Cette rubrique présente un exemple de client Go utilisant le protocole AMQP pour se connecter à Alibaba Cloud IoT Platform et recevoir des messages via un abonnement côté serveur.
Prérequis
Vous disposez d'un ID de groupe de consommateurs et vous êtes abonné aux messages de la rubrique requise.
Gérer les groupes de consommateurs AMQP : utilisez le groupe de consommateurs par défaut (DEFAULT_GROUP) dans IoT Platform ou créez-en un nouveau.
Configurer un abonnement côté serveur AMQP : abonnez-vous aux messages de la rubrique requise en utilisant un groupe de consommateurs.
Préparer l'environnement de développement
Cet exemple utilise Go 1.12.7.
Télécharger le SDK
Importez le SDK AMQP pour Go à l'aide de la commande suivante.
import "pack.ag/amqp"
Pour plus d'informations sur l'utilisation du SDK, consultez package amqp.
Exemple de code
package main
import (
"os"
"context"
"crypto/hmac"
"crypto/sha1"
"encoding/base64"
"fmt"
"pack.ag/amqp"
"time"
)
// For parameter descriptions, see the AMQP client connection guide.
const consumerGroupId = "${YourConsumerGroupId}"
const clientId = "${YourClientId}"
// iotInstanceId: The instance ID.
const iotInstanceId = "${YourIotInstanceId}"
// The endpoint. For more information, see the AMQP client connection guide.
const host = "${YourHost}"
func main() {
// If you leak the project code, your AccessKey may be exposed. This compromises the security of all resources in your account.
// The following code provides an example on how to use an environment variable to obtain the AccessKey. This method is for reference only.
accessKey := os.Getenv("ALIBABA_CLOUD_ACCESS_KEY_ID")
accessSecret := os.Getenv("ALIBABA_CLOUD_ACCESS_KEY_SECRET")
address := "amqps://" + host + ":5671"
timestamp := time.Now().Nanosecond() / 1000000
// For information about how to construct the user name, see the AMQP client connection guide.
userName := fmt.Sprintf("%s|authMode=aksign,signMethod=Hmacsha1,consumerGroupId=%s,authId=%s,iotInstanceId=%s,timestamp=%d|",
clientId, consumerGroupId, accessKey, iotInstanceId, timestamp)
stringToSign := fmt.Sprintf("authId=%s×tamp=%d", accessKey, timestamp)
hmacKey := hmac.New(sha1.New, []byte(accessSecret))
hmacKey.Write([]byte(stringToSign))
// Calculate the signature. For information about how to construct the password, see the AMQP client connection guide.
password := base64.StdEncoding.EncodeToString(hmacKey.Sum(nil))
amqpManager := &AmqpManager{
address:address,
userName:userName,
password:password,
}
// To accept messages or cancel operations, derive the context from context.Background().
ctx := context.Background()
amqpManager.startReceiveMessage(ctx)
}
// Business function. This is a custom implementation. The function is executed asynchronously. Consider the resource consumption of your system.
func (am *AmqpManager) processMessage(message *amqp.Message) {
fmt.Println("data received:", string(message.GetData()), " properties:", message.ApplicationProperties)
}
type AmqpManager struct {
address string
userName string
password string
client *amqp.Client
session *amqp.Session
receiver *amqp.Receiver
}
func (am *AmqpManager) startReceiveMessage(ctx context.Context) {
childCtx, _ := context.WithCancel(ctx)
err := am.generateReceiverWithRetry(childCtx)
if nil != err {
return
}
defer func() {
am.receiver.Close(childCtx)
am.session.Close(childCtx)
am.client.Close()
}()
for {
// Blocks to receive messages. If the context is background, the process is not interrupted.
message, err := am.receiver.Receive(ctx)
if nil == err {
go am.processMessage(message)
message.Accept()
} else {
fmt.Println("amqp receive data error:", err)
// If the operation is actively canceled, exit the program.
select {
case <- childCtx.Done(): return
default:
}
// If the operation is not actively canceled, re-establish the connection.
err := am.generateReceiverWithRetry(childCtx)
if nil != err {
return
}
}
}
}
func (am *AmqpManager) generateReceiverWithRetry(ctx context.Context) error {
// Back off and reconnect. Start with 10 ms and double the interval up to 20s.
duration := 10 * time.Millisecond
maxDuration := 20000 * time.Millisecond
times := 1
// In case of an exception, back off and reconnect.
for {
select {
case <- ctx.Done(): return amqp.ErrConnClosed
default:
}
err := am.generateReceiver()
if nil != err {
time.Sleep(duration)
if duration < maxDuration {
duration *= 2
}
fmt.Println("amqp connect retry,times:", times, ",duration:", duration)
times ++
} else {
fmt.Println("amqp connect init success")
return nil
}
}
}
// The connection and session states cannot be determined because the package is not visible. Restart the connection to retrieve the states.
func (am *AmqpManager) generateReceiver() error {
if am.session != nil {
receiver, err := am.session.NewReceiver(
amqp.LinkSourceAddress("/queue-name"),
amqp.LinkCredit(20),
)
// If a network disconnection occurs, the connection closes and session creation fails.
// If the connection remains open, the session is created successfully.
if err == nil {
am.receiver = receiver
return nil
}
}
// Clean up the previous connection.
if am.client != nil {
am.client.Close()
}
client, err := amqp.Dial(am.address, amqp.ConnSASLPlain(am.userName, am.password), )
if err != nil {
return err
}
am.client = client
session, err := client.NewSession()
if err != nil {
return err
}
am.session = session
receiver, err := am.session.NewReceiver(
amqp.LinkSourceAddress("/queue-name"),
amqp.LinkCredit(20),
)
if err != nil {
return err
}
am.receiver = receiver
return nil
}
Configurez les paramètres du code précédent comme indiqué dans le tableau ci-dessous. Pour plus d'informations, consultez Connecter un client AMQP à IoT Platform.
Spécifiez des valeurs de paramètres valides. Sinon, le client AMQP ne pourra pas se connecter à IoT Platform.
|
Paramètre |
Description |
|
accessKey |
Connectez-vous à la console IoT Platform, placez le curseur sur votre photo de profil, puis cliquez sur AccessKey Management pour obtenir l'ID AccessKey et le secret AccessKey. Remarque
Si vous utilisez un utilisateur Resource Access Management (RAM), accordez-lui l'autorisation AliyunIOTFullAccess. Cette autorisation est requise pour gérer IoT Platform. Sans elle, la connexion échoue. Pour savoir comment accorder des autorisations, consultez Accès utilisateur RAM. |
|
accessSecret |
|
|
consumerGroupId |
L'ID du groupe de consommateurs dans l'instance IoT Platform. Connectez-vous à la console IoT Platform. Dans l'instance correspondante, accédez à pour afficher l'ID de votre groupe de consommateurs. |
|
iotInstanceId |
L'ID de l'instance. Vous pouvez afficher l'ID de l'instance actuelle sur la page Instance Overview de la console IoT Platform.
|
|
clientId |
L'ID client. Vous devez définir cet ID. Sa longueur peut atteindre 64 caractères. Nous vous recommandons d'utiliser un identifiant unique, tel que l'UUID, l'adresse MAC ou l'adresse IP du serveur où se trouve votre client AMQP. Une fois le client AMQP connecté et démarré, connectez-vous à la console IoT Platform. Sur l'onglet Consumer Groups de la page de l'instance, cliquez sur View à côté du groupe de consommateurs. La page Consumer Group Details affiche ce paramètre. Cela vous aide à identifier différents clients. |
|
host |
Le point de terminaison AMQP. Pour connaître le point de terminaison AMQP correspondant à |
Résultats d'exécution de l'exemple
-
Succès : Un message de journal similaire au suivant s'affiche. Cela indique que le client AMQP s'est connecté avec succès à IoT Platform et a reçu des messages de l'appareil.
amqp connect init success data received: {"deviceType":"CustomCategory","iotId":"xxx","requestId":"1613726251726","checkFailedData":0,"productKey":"xxx","gmtCreate":"1613726121717,"deviceName":"xxx","items":{"Temperature":{"value":24,"time":1613726121715},"Humidity":{"value":19,"time":1613726121715}}} properties: map[generateTime:1613726121721 messageId:xxx xxx s:1 topic: /xxx/thing/event/property/post] data received: {"deviceType":"CustomCategory","iotId":"xxx","requestId":"1613725651726","checkFailedData":0,"productKey":"xxx","gmtCreate":"1613725521715,"deviceName":"xxx","items":{"Temperature":{"value":28,"time":1613725521712},"Humidity":{"value":19,"time":1613725521712}}} properties: map[generateTime:1613725521719 messageId:1362689721104473600 qos:1 topic: /xxx/thing/event/property/post] -
Échec : Un message de journal similaire au suivant s'affiche, indiquant que le client AMQP n'a pas réussi à se connecter à IoT Platform.
Utilisez le journal des erreurs pour vérifier votre code et vos paramètres réseau. Résolvez le problème et exécutez à nouveau le code.
amqp connect retry,times: 1 ,duration: 20ms amqp connect retry,times: 2 ,duration: 40ms amqp connect retry,times: 3 ,duration: 80ms amqp connect retry,times: 4 ,duration: 160ms amqp connect retry,times: 5 ,duration: 320ms amqp connect retry,times: 6 ,duration: 640ms amqp connect retry,times: 7 ,duration: 1.28s amqp connect retry,times: 8 ,duration: 2.56s amqp connect retry,times: 9 ,duration: 5.12s amqp connect retry,times: 10 ,duration: 10.24s amqp connect retry,times: 11 ,duration: 20.48s
Références
Pour plus d'informations sur les codes d'erreur des messages d'abonnement côté serveur, consultez Codes d'erreur liés aux messages.