Tous les produits
Search
Centre de documentation

ApsaraMQ for RocketMQ:Sample code

Dernière mise à jour :Aug 09, 2026

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 :

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()
Remarque

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

Voir aussi