Envoyez et recevez des messages avec le SDK ApsaraMQ for RocketMQ 5.x pour Python. Chaque exemple traite d'un type de message spécifique et est prêt à l'emploi après remplacement des valeurs d'espace réservé par les détails de votre instance.
Prérequis
Avant de commencer, vérifiez que vous disposez des éléments suivants :
Le SDK RocketMQ pour Python installé. Pour plus de détails sur les versions, consultez le Guide des versions
Une instance ApsaraMQ for RocketMQ 5.x avec un endpoint. Pour les instructions de configuration, consultez la section Préparation de l'environnement
Un topic et un consumer group créés sur l'instance
Paramètres
Remplacez les espaces réservés suivants dans chaque exemple de code par vos valeurs réelles.
| Paramètre | Valeur d'exemple | Description |
|---|---|---|
<your-endpoint> |
rmq-cn-xxx.{regionId}.rmq.aliyuncs.com:8080 |
L'endpoint de votre instance ApsaraMQ for RocketMQ. Pour obtenir l'endpoint, consultez la section Obtenir l'endpoint de l'instance. Utilisez l'endpoint public pour un accès Internet ou l'endpoint VPC pour un accès VPC. |
<your-instance-id> |
rmq-cn-xxx |
L'ID de votre instance ApsaraMQ for RocketMQ. Requis uniquement pour l'accès au réseau public sur les instances serverless. |
<your-topic> |
normal_test |
Le topic vers lequel les messages sont envoyés ou depuis lequel ils sont consommés. Créez le topic au préalable. |
<your-consumer-group> |
GID_test |
Le consumer group utilisé pour consommer les messages. Créez le consumer group au préalable. |
<your-ak> |
1XVg0hzgKm****** |
Le nom d'utilisateur de votre instance. Requis pour l'accès Internet. Pour l'accès VPC, requis uniquement si l'instance est serverless et que l'authentification sans mot de passe dans les VPC est désactivée. Consultez la section Obtenir le nom d'utilisateur et le mot de passe de l'instance. |
<your-sk> |
ijSt8rEc45****** |
Le mot de passe de votre instance. Mêmes conditions d'accès que pour le nom d'utilisateur. |
Envoyer des messages normaux
Les messages normaux n'ont pas de sémantique de livraison spéciale et constituent le type de message le plus courant.
Envoi synchrone
La méthode send() bloque jusqu'à ce que le broker accuse réception du message.
from rocketmq import ClientConfiguration, Credentials, Message, Producer
if __name__ == '__main__':
# Instance endpoint
endpoints = "<your-endpoint>"
# Authenticate with your instance username and password
credentials = Credentials("<your-ak>", "<your-sk>")
config = ClientConfiguration(endpoints, credentials)
topic = "<your-topic>"
producer = Producer(config, (topic,))
try:
producer.startup()
try:
msg = Message()
msg.topic = topic
msg.body = "hello, rocketmq.".encode('utf-8')
# Tag: secondary classifier of messages besides topic
msg.tag = "tag"
# Keys: alternative way to identify messages besides message ID
msg.keys = "keys"
# Custom user property
msg.add_property("key", "value")
for i in range(10):
res = producer.send(msg)
print(f"Send message success. {res}")
producer.shutdown()
except Exception as e:
print(f"Send message failed: {e}")
producer.shutdown()
except Exception as e:
print(f"Producer startup failed: {e}")
producer.shutdown()
Envoi asynchrone
La méthode send_async() renvoie immédiatement le contrôle et invoque un callback lorsque le broker répond. Utilisez l'envoi asynchrone pour un débit plus élevé lorsque vous n'avez pas besoin d'attendre chaque accusé de réception.
from rocketmq import ClientConfiguration, Credentials, Message, Producer
def handle_send_result(result_future):
try:
# Avoid time-consuming logic in the callback; offload to another thread if needed
res = result_future.result()
print(f"Send message success. {res}")
except Exception as exception:
print(f"Send message failed: {exception}")
if __name__ == '__main__':
endpoints = "<your-endpoint>"
credentials = Credentials("<your-ak>", "<your-sk>")
config = ClientConfiguration(endpoints, credentials)
topic = "<your-topic>"
producer = Producer(config, (topic,))
try:
producer.startup()
try:
for i in range(10):
msg = Message()
msg.topic = topic
msg.body = "hello, rocketmq.".encode('utf-8')
msg.tag = "tag"
msg.keys = "keys"
msg.add_property("send", "async")
send_result_future = producer.send_async(msg)
send_result_future.add_done_callback(handle_send_result)
except Exception as e:
print(f"Send message failed: {e}")
except Exception as e:
print(f"Producer startup failed: {e}")
input("Press Enter to stop the application.")
producer.shutdown()
Envoyer des messages ordonnés
Les messages ordonnés sont livrés dans l'ordre d'envoi au sein du même groupe de messages. Définissez la propriété message_group pour regrouper les messages associés.
from rocketmq import ClientConfiguration, Credentials, Message, Producer
if __name__ == '__main__':
endpoints = "<your-endpoint>"
credentials = Credentials("<your-ak>", "<your-sk>")
config = ClientConfiguration(endpoints, credentials)
topic = "<your-topic>"
producer = Producer(config, (topic,))
try:
producer.startup()
try:
msg = Message()
msg.topic = topic
msg.body = "hello, rocketmq.".encode('utf-8')
msg.tag = "rocketmq-send-fifo-message"
msg.keys = "keys"
# Messages with the same message_group are delivered in order
msg.message_group = "fifo-group"
for i in range(10):
res = producer.send(msg)
print(f"Send message success. {res}")
producer.shutdown()
except Exception as e:
print(f"Send message failed: {e}")
producer.shutdown()
except Exception as e:
print(f"Producer startup failed: {e}")
producer.shutdown()
Le code source pour les messages ordonnés est disponible à l'adresse fifo_producer_example.py.
Envoyer des messages planifiés et différés
Les messages planifiés et différés sont livrés après un horodatage spécifié. Définissez delivery_timestamp sur une heure Unix en secondes.
import time
from rocketmq import ClientConfiguration, Credentials, Message, Producer
if __name__ == '__main__':
endpoints = "<your-endpoint>"
credentials = Credentials("<your-ak>", "<your-sk>")
config = ClientConfiguration(endpoints, credentials)
topic = "<your-topic>"
producer = Producer(config, (topic,))
try:
producer.startup()
try:
msg = Message()
msg.topic = topic
msg.body = "hello, rocketmq.".encode('utf-8')
msg.tag = "rocketmq-send-delay-message"
# Deliver the message 10 seconds from now
msg.delivery_timestamp = int(time.time()) + 10
res = producer.send(msg)
print(f"Send message success. {res}")
producer.shutdown()
except Exception as e:
print(f"Send message failed: {e}")
producer.shutdown()
except Exception as e:
print(f"Producer startup failed: {e}")
producer.shutdown()
Envoyer des messages transactionnels
Les messages transactionnels prennent en charge la validation en deux phases. Le broker contacte le producteur pour confirmer ou annuler les messages non validés via un TransactionChecker.
from rocketmq import (ClientConfiguration, Credentials, Message, Producer,
TransactionChecker, TransactionResolution)
class OrderTransactionChecker(TransactionChecker):
"""Check whether the local transaction succeeded and return COMMIT or ROLLBACK."""
def check(self, message: Message) -> TransactionResolution:
print(f"Transaction check for message: {message}. Committing.")
return TransactionResolution.COMMIT
if __name__ == '__main__':
endpoints = "<your-endpoint>"
credentials = Credentials("<your-ak>", "<your-sk>")
config = ClientConfiguration(endpoints, credentials)
topic = "<your-topic>"
# Pass the transaction checker to the producer
producer = Producer(config, (topic,), checker=OrderTransactionChecker())
try:
producer.startup()
except Exception as e:
print(f"Producer startup failed: {e}")
try:
# Begin a half-message transaction
transaction = producer.begin_transaction()
msg = Message()
msg.topic = topic
msg.body = "hello, rocketmq.".encode('utf-8')
msg.tag = "rocketmq-send-transaction-message"
res = producer.send(msg, transaction)
print(f"Send half message success. {res}")
# Option 1: Commit from client side
transaction.commit()
print(f"Committed message: {transaction.message_id}")
# Option 2: Roll back instead
# transaction.rollback()
# print(f"Rolled back message: {transaction.message_id}")
producer.shutdown()
except Exception as e:
print(f"Transaction send failed: {e}")
producer.shutdown()
Envoyer des messages légers (Lite)
Les messages légers utilisent un topic parent avec des sous-topics (lite topics) pour catégoriser les messages avec une granularité plus fine. Définissez lite_topic sur le message pour spécifier le sous-topic.
from rocketmq import ClientConfiguration, Credentials, Message, Producer
if __name__ == '__main__':
endpoints = "<your-endpoint>"
credentials = Credentials("<your-ak>", "<your-sk>")
config = ClientConfiguration(endpoints, credentials)
# Use the parent topic
topic = "<your-topic>"
producer = Producer(config, (topic,))
try:
producer.startup()
try:
msg = Message()
msg.topic = topic
msg.body = "hello, rocketmq.".encode('utf-8')
for i in range(10):
# Route each message to a lite topic under the parent topic
msg.lite_topic = f"lite-test-{i}"
res = producer.send(msg)
print(f"Send message success. {res}")
producer.shutdown()
except Exception as e:
print(f"Send message failed: {e}")
producer.shutdown()
except Exception as e:
print(f"Producer startup failed: {e}")
producer.shutdown()
Consommer des messages avec un simple consumer
Un simple consumer extrait les messages du broker via un long polling. Après avoir traité chaque message, appelez ack() pour l'acquitter.
from rocketmq import (ClientConfiguration, Credentials, FilterExpression,
SimpleConsumer)
if __name__ == '__main__':
endpoints = "<your-endpoint>"
credentials = Credentials("<your-ak>", "<your-sk>")
config = ClientConfiguration(endpoints, credentials)
topic = "<your-topic>"
consumer_group = "<your-consumer-group>"
# Use a singleton consumer instance in production; avoid creating multiple consumers
simple_consumer = SimpleConsumer(config, consumer_group, {topic: FilterExpression()})
try:
simple_consumer.startup()
try:
# Optional: filter by tag
# simple_consumer.subscribe(topic, FilterExpression("your-tag"))
while True:
try:
# Long poll for up to 32 messages; each message is invisible
# for 15 seconds after receipt
messages = simple_consumer.receive(32, 15)
if messages is not None:
for msg in messages:
simple_consumer.ack(msg)
print(f"Acknowledged message: [{msg.message_id}]")
except Exception as e:
print(f"Receive or ack failed: {e}")
except Exception as e:
print(f"Consumer error: {e}")
simple_consumer.shutdown()
except Exception as e:
print(f"Consumer startup failed: {e}")
simple_consumer.shutdown()
Consommer des messages légers avec un push consumer
Un push consumer reçoit automatiquement les messages via un listener enregistré. Pour l'exemple de push consumer pour les messages légers, consultez lite_push_consumer_example.py.
Accès au réseau public pour les instances serverless
Pour accéder à une instance ApsaraMQ for RocketMQ serverless via le réseau public, transmettez l'ID de l'instance en tant que troisième argument à ClientConfiguration :
config = ClientConfiguration(endpoints, credentials, "<your-instance-id>")
Remplacez <your-instance-id> par l'ID réel de votre instance, par exemple rmq-cn-xxx.
Fichiers source GitHub
Le tableau suivant associe chaque type de message à son fichier source GitHub.
| Type de message | Envoi | Consommation |
|---|---|---|
| Messages normaux | Synchrone : normal_producer_example.py, Asynchrone : async_producer_example.py | simple_consumer_example.py |
| Messages ordonnés | fifo_producer_example.py | - |
| Messages planifiés et différés | delay_producer_example.py | - |
| Messages transactionnels | transaction_producer_example.py | - |
| Messages légers | lite_producer_example.py | lite_push_consumer_example.py |