Utilisez le SDK pour Python 2.7 afin de connecter un client AMQP à Alibaba Cloud IoT Platform et de 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 un groupe de consommateurs.
Configurer un abonnement côté serveur AMQP : abonnez-vous aux messages de la rubrique requise à l'aide d'un groupe de consommateurs.
Environnement de développement
Cet exemple utilise Python 2.7.
Télécharger le SDK
Nous recommandons la bibliothèque Apache Qpid Proton 0.29.0, qui encapsule l'API Python. Pour télécharger la bibliothèque et consulter les instructions, accédez à Qpid Proton 0.29.0.
Installez Qpid Proton. Pour plus d'informations, consultez Installation de Qpid Proton.
Après avoir installé Qpid Proton, exécutez la commande Python suivante pour vérifier que la bibliothèque SSL est disponible :
import proton;print('%s' % 'SSL present' if proton.SSL.present() else 'SSL NOT AVAILABLE')
Exemple de code
# encoding=utf-8
import sys
import logging
import time
from proton.handlers import MessagingHandler
from proton.reactor import Container
import hashlib
import hmac
import base64
import os
reload(sys)
sys.setdefaultencoding('utf-8')
logging.basicConfig(level=logging.INFO, format='%(asctime)s - %(name)s - %(levelname)s - %(message)s')
logger = logging.getLogger(__name__)
console_handler = logging.StreamHandler(sys.stdout)
def current_time_millis():
return str(int(round(time.time() * 1000)))
def do_sign(secret, sign_content):
m = hmac.new(secret, sign_content, digestmod=hashlib.sha1)
return base64.b64encode(m.digest())
class AmqpClient(MessagingHandler):
def __init__(self):
super(AmqpClient, self).__init__()
def on_start(self, event):
# The endpoint used to connect to IoT Platform. For details, see the AMQP connection guide.
url = "amqps://${YourHost}:5671"
# Hard-coding an AccessKey pair in your project code poses a security risk. If the code is leaked, your AccessKey pair will be exposed, compromising all resources in your account. The following code retrieves the AccessKey pair from environment variables as a best practice. This is for reference only.
accessKey = os.environ['ALIBABA_CLOUD_ACCESS_KEY_ID']
accessSecret = os.environ['ALIBABA_CLOUD_ACCESS_KEY_SECRET']
consumerGroupId = "${YourConsumerGroupId}"
clientId = "${YourClientId}"
# iotInstanceId: The ID of the IoT Platform instance.
iotInstanceId = "${YourIotInstanceId}"
# The signature algorithm. Valid values: hmacmd5, hmacsha1, and hmacsha256.
signMethod = "hmacsha1"
timestamp = current_time_millis()
# For instructions on how to construct the userName parameter, see the AMQP connection guide.
userName = clientId + "|authMode=aksign" + ",signMethod=" + signMethod \
+ ",timestamp=" + timestamp + ",authId=" + accessKey \
+ ",iotInstanceId=" + iotInstanceId + ",consumerGroupId=" + consumerGroupId + "|"
signContent = "authId=" + accessKey + "×tamp=" + timestamp
# Calculate the signature. For instructions on how to construct the password, see the AMQP connection guide.
passWord = do_sign(accessSecret.encode("utf-8"), signContent.encode("utf-8"))
conn = event.container.connect(url, user=userName, password=passWord, heartbeat=60)
self.receiver = event.container.create_receiver(conn)
# Called when the connection is successfully established.
def on_connection_opened(self, event):
logger.info("Connection established, remoteUrl: %s", event.connection.hostname)
# Called when the connection is closed.
def on_connection_closed(self, event):
logger.info("Connection closed: %s", self)
# Called when the remote peer closes the connection due to an error.
def on_connection_error(self, event):
logger.info("Connection error")
# Called when an AMQP connection error occurs, such as an authentication or socket error.
def on_transport_error(self, event):
if event.transport.condition:
if event.transport.condition.info:
logger.error("%s: %s: %s" % (
event.transport.condition.name, event.transport.condition.description,
event.transport.condition.info))
else:
logger.error("%s: %s" % (event.transport.condition.name, event.transport.condition.description))
else:
logging.error("Unspecified transport error")
# Called when a message is received.
def on_message(self, event):
message = event.message
content = message.body.decode('utf-8')
topic = message.properties.get("topic")
message_id = message.properties.get("messageId")
print("receive message: message_id=%s, topic=%s, content=%s" % (message_id, topic, content))
event.receiver.flow(1)
Container(AmqpClient()).run()
Configurez les paramètres du code précédent comme décrit dans le tableau suivant. Pour plus d'informations, consultez Connecter un client AMQP à IoT Platform.
Spécifiez des valeurs de paramètre valides. Sinon, le client AMQP ne parviendra pas à se connecter à IoT Platform.
|
Paramètre |
Description |
|
url |
Le point de terminaison utilisé par le client AMQP pour se connecter à IoT Platform. Format : Pour plus d'informations sur le point de terminaison que vous pouvez spécifier pour la variable |
|
accessKey |
Connectez-vous à la console IoT Platform, survolez l'image de profil dans le coin supérieur droit, puis cliquez sur AccessKey Management pour obtenir l'ID AccessKey et le secret AccessKey. Remarque
Si vous utilisez un utilisateur Resource Access Management (RAM), vous devez joindre la politique AliyunIOTFullAccess à l'utilisateur RAM pour lui accorder les autorisations de gestion des ressources IoT Platform. Sinon, la connexion échoue. Pour plus d'informations, consultez Accéder à IoT Platform en tant qu'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 votre ID de groupe de consommateurs. |
|
iotInstanceId |
L'ID de l'instance IoT Platform. Vous pouvez afficher l'ID de l'instance dans l'onglet Overview de la console IoT Platform.
|
|
clientId |
L'ID client. Vous devez définir cet ID. L'ID peut comporter jusqu'à 64 caractères. Nous vous recommandons d'utiliser un identifiant unique, tel que l'UUID, l'adresse MAC ou l'adresse IP du serveur sur lequel 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. |
Résultats d'exemple
-
Succès : Le message de journal suivant indique que le client AMQP s'est connecté à IoT Platform et reçoit des messages.
202x-xx-xx xx:xx:xx,490 - __main__ - INFO - Connection established, remoteUrl: xxx.amqp.iothub.aliyuncs.com receive message: message_id=13xxx, topic=/xxx/thing/event/property/post, content={"deviceType":"CustomCategory","iotId":"xxx","requestId":"1xxx"} receive message: message_id=13xxx, topic=/xxx/thing/event/property/post, content={"deviceType":"CustomCategory","iotId":"xxx","requestId":"1xxx"}Paramètre
Exemple
Description
message_id
2****7
L'ID du message.
topic
//*/thing/event/property/post
La rubrique utilisée pour soumettre les propriétés de l'appareil.
content
{"deviceType":"CustomCategory","iotId":"qPi","requestId":"161","checkFailedData":{},"productKey":"g4","gmtCreate":1613635594038,"deviceName":"de","items":{"Temperature":{"value":24,"time":1613635594036},"Humidity":{"value":26,"time":1613635594036}}}
Le contenu du message.
-
Si des informations similaires à la sortie suivante s'affichent, le client AMQP ne parvient pas à se connecter à IoT Platform.
Vérifiez votre code et vos paramètres réseau à l'aide du journal des erreurs. Résolvez le problème et réexécutez le code.
2021-02-18 16:19:40,993 - proton - ERROR - Couldn't connect: 10053 2021-02-18 16:19:40,996 - __main__ - ERROR - proton.pythonio: Connection error: 10053 2021-02-18 16:19:40,996 - proton - INFO - Disconnected, reconnecting... Traceback (most recent call last): File "xxx", line 87, in <module> Container(AmqpClient()).run() File "xxx", line 181, in run while self.process(): pass File "xxx", line 240, in process event.dispatch(self._global_handler) File "xxx", line 135, in dispatch _dispatch(handler, type.method, self) File "xxx", line 117, in _dispatch handler.on_unhandled(method, *args) File "xxx", line 665, in on_unhandled event.dispatch(self.base) File "xxx", line 135, in dispatch _dispatch(handler, type.method, self) File "xxx", line 115, in _dispatch m(*args) File "xxx", line 892, in on_connection_bound addrs = socket.getaddrinfo(host, port, socket.AF_UNSPEC, socket.SOCK_STREAM) socket.gaierror: [Errno 11001] getaddrinfo failed
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.