Tous les produits
Search
Centre de documentation

ApsaraMQ for RocketMQ:Send and receive ordered messages

Dernière mise à jour :Aug 09, 2026

Les messages ordonnés, également appelés messages FIFO (First In, First Out), sont publiés et consommés dans un ordre strict. Cette rubrique fournit des exemples de code C++ utilisant le SDK client TCP pour C++ de l'édition Community afin d'envoyer et de recevoir des messages ordonnés.

Fonctionnement

ApsaraMQ for RocketMQ prend en charge deux types de messages ordonnés :

Type Étendue de l'ordre Description
Globalement ordonné Rubrique entière Tous les messages d'une rubrique sont publiés et consommés selon l'ordre FIFO.
Ordonné par partition Par partition Les messages sont distribués aux partitions via une clé de sharding. Les messages au sein de chaque partition sont publiés et consommés selon l'ordre FIFO.

Une clé de sharding est un champ clé qui identifie différentes partitions pour les messages ordonnés. Elle diffère de la clé d'un message standard.

Pour plus d'informations, consultez la section Messages ordonnés.

Prérequis

Avant de commencer, assurez-vous d'avoir :

Envoyer des messages ordonnés

Important

Le broker ApsaraMQ for RocketMQ détermine l'ordre des messages en fonction de l'ordre d'envoi par un producteur unique ou un thread unique. Si l'expéditeur utilise plusieurs producteurs ou threads pour envoyer des messages simultanément, le broker trie les messages par heure d'arrivée, ce qui peut différer de l'ordre métier prévu.

  1. Enregistrez le code suivant sous le nom OrderProducerDemo.cpp, puis compilez-le :

    g++ -o order_producer_demo -std=c++11 -lz -lrocketmq OrderProducerDemo.cpp

  2. Exécutez le binaire compilé. Le code envoie 32 messages ordonnés, en acheminant chacun vers une file d'attente basée sur la clé de partition :

    #include <iostream>
    #include <chrono>
    #include <thread>
    #include "DefaultMQProducer.h"
    
    using namespace std;
    using namespace rocketmq;
    
    // Custom queue selector: routes messages to a queue based on custom
    // partitioning logic. In this example, the modulo of orderId is used.
    class ExampleSelectMessageQueueByHash : public MessageQueueSelector {
    public:
        MQMessageQueue select(const std::vector<MQMessageQueue> &mqs, const MQMessage &msg, void *arg) {
            int orderId = *static_cast<int *>(arg);
            int index = orderId % mqs.size();
            return mqs[0];
        }
    };
    
    int main() {
        std::cout << "=======Before sending messages=======" << std::endl;
    
        // Replace with the Group ID that you created in the ApsaraMQ for RocketMQ console.
        DefaultMQProducer producer("GID_XXXXXXXX");
    
        // Replace with the TCP endpoint from the Instance Details page
        // in the ApsaraMQ for RocketMQ console.
        producer.setNamesrvAddr("http://MQ_INST_XXXXXXXXXX.mq-internet-access.mq-internet.aliyuncs.com:80");
    
        // Set credentials from environment variables.
        // ALIBABA_CLOUD_ACCESS_KEY_ID: your AccessKey ID
        // ALIBABA_CLOUD_ACCESS_KEY_SECRET: your AccessKey secret
        // "ALIYUN": the user channel (default value)
        producer.setSessionCredentials(
            getenv("ALIBABA_CLOUD_ACCESS_KEY_ID"),
            getenv("ALIBABA_CLOUD_ACCESS_KEY_SECRET"),
            "ALIYUN"
        );
    
        producer.start();
    
        auto start = std::chrono::system_clock::now();
        int count = 32;
        int retryTimes = 1;
    
        // Create the queue selector for ordered message routing.
        ExampleSelectMessageQueueByHash *pSelector = new ExampleSelectMessageQueueByHash();
    
        for (int i = 0; i < count; ++i) {
            // Replace "YOUR ORDERLY TOPIC" with the topic that you created
            // in the ApsaraMQ for RocketMQ console.
            MQMessage msg("YOUR ORDERLY TOPIC", "HiTAG", "Hello,CPP SDK, Orderly Message.");
    
            try {
                // Send the message to the queue determined by the selector.
                // The fourth argument (&i) serves as the partition key.
                SendResult sendResult = producer.send(msg, pSelector, &i, 1, false);
                std::cout << "SendResult:" << sendResult.getSendStatus()
                          << ", Message ID: " << sendResult.getMsgId()
                          << "MessageQueue:" << sendResult.getMessageQueue().toString()
                          << std::endl;
                this_thread::sleep_for(chrono::seconds(1));
            } catch (MQException e) {
                std::cout << "ErrorCode: " << e.GetError()
                          << " Exception:" << e.what() << std::endl;
            }
        }
    
        auto interval = std::chrono::system_clock::now() - start;
        std::cout << "Send " << count << " messages OK, costs "
                  << std::chrono::duration_cast<std::chrono::milliseconds>(interval).count()
                  << "ms" << std::endl;
    
        producer.shutdown();
        std::cout << "=======After sending messages=======" << std::endl;
        return 0;
    }

Remplacez les espaces réservés suivants par des valeurs réelles :

Espace réservé Description Exemple
GID_XXXXXXXX L'ID de groupe depuis la console ApsaraMQ for RocketMQ GID_OrderGroup
MQ_INST_XXXXXXXXXX.mq-internet-access.mq-internet.aliyuncs.com:80 Le point de terminaison TCP depuis la page Instance Details MQ_INST_abc123.mq-internet-access.mq-internet.aliyuncs.com:80
YOUR ORDERLY TOPIC La rubrique créée dans la console ApsaraMQ for RocketMQ OrderTopic

Consommer des messages ordonnés

Le consommateur utilise MessageListenerOrderly pour traiter les messages dans un ordre FIFO strict au sein de chaque file d'attente. L'exemple de code suivant utilise l'édition commerciale du SDK client TCP pour C++.

  1. Enregistrez le code suivant sous le nom OrderConsumerDemo.cpp, puis compilez-le :

    g++ -o order_consumer_demo -std=c++11 -lz -lrocketmq OrderConsumerDemo.cpp

  2. Exécutez le binaire compilé. Le code s'abonne à une rubrique et consomme les messages ordonnés de manière séquentielle :

    #include <iostream>
    #include <thread>
    #include "DefaultMQPushConsumer.h"
    
    using namespace rocketmq;
    
    // Listener that processes messages in the order they were stored.
    class ExampleOrderlyMessageListener : public MessageListenerOrderly {
    public:
        ConsumeStatus consumeMessage(const std::vector<MQMessageExt> &msgs) {
            for (auto item = msgs.begin(); item != msgs.end(); item++) {
                std::cout << "Received Message Topic:" << item->getTopic()
                          << ", MsgId:" << item->getMsgId() << std::endl;
            }
            return CONSUME_SUCCESS;
        }
    };
    
    int main(int argc, char *argv[]) {
        std::cout << "=======Before consuming messages=======" << std::endl;
    
        // Replace with the Group ID from the ApsaraMQ for RocketMQ console.
        DefaultMQPushConsumer *consumer = new DefaultMQPushConsumer("GID_XXXXXXXXXXX");
    
        // Replace with the TCP endpoint from the Instance Details page
        // in the ApsaraMQ for RocketMQ console.
        consumer->setNamesrvAddr(
            "http://MQ_INST_XXXXXXXXXXXXXX.mq-internet-access.mq-internet.aliyuncs.com:80"
        );
    
        // Set credentials from environment variables.
        // ALIBABA_CLOUD_ACCESS_KEY_ID: your AccessKey ID
        // ALIBABA_CLOUD_ACCESS_KEY_SECRET: your AccessKey secret
        // "ALIYUN": the user channel (default value)
        consumer->setSessionCredentials(
            getenv("ALIBABA_CLOUD_ACCESS_KEY_ID"),
            getenv("ALIBABA_CLOUD_ACCESS_KEY_SECRET"),
            "ALIYUN"
        );
    
        auto start = std::chrono::system_clock::now();
    
        // Register the orderly message listener.
        ExampleOrderlyMessageListener *messageListener = new ExampleOrderlyMessageListener();
    
        // Subscribe to the topic. The "*" expression matches all tags.
        consumer->subscribe("YOUR ORDERLY TOPIC", "*");
        consumer->registerMessageListener(messageListener);
    
        // Start the consumer.
        // Before starting:
        //   - Make sure that the subscription is configured.
        //   - Make sure that all consumers in the same group use consistent subscriptions.
        consumer->start();
    
        // Keep the main thread alive for 60 seconds to receive messages.
        std::this_thread::sleep_for(std::chrono::seconds(60));
        consumer->shutdown();
    
        std::cout << "=======After consuming messages======" << std::endl;
        return 0;
    }

Remplacez les espaces réservés suivants par des valeurs réelles :

Espace réservé Description Exemple
GID_XXXXXXXXXXX L'ID de groupe depuis la console ApsaraMQ for RocketMQ GID_OrderConsumerGroup
MQ_INST_XXXXXXXXXXXXXX.mq-internet-access.mq-internet.aliyuncs.com:80 Le point de terminaison TCP depuis la page Instance Details MQ_INST_abc123.mq-internet-access.mq-internet.aliyuncs.com:80
YOUR ORDERLY TOPIC La rubrique créée dans la console ApsaraMQ for RocketMQ OrderTopic

Étapes suivantes

  • Consultez la section Messages ordonnés pour comprendre en détail le fonctionnement de l'ordonnancement global et par partition.

  • Explorez les autres types de messages pris en charge par le SDK client TCP pour C++.