Tous les produits
Search
Centre de documentation

ApsaraMQ for RocketMQ:Subscribe to messages

Dernière mise à jour :Aug 09, 2026

Abonnez-vous aux messages d’ApsaraMQ for RocketMQ à l’aide du SDK client TCP pour C/C++. Cette rubrique présente les deux modes d’abonnement, propose une implémentation complète du consommateur et détaille la gestion de son cycle de vie.

Modes d’abonnement

ApsaraMQ for RocketMQ prend en charge deux modes d’abonnement qui déterminent la distribution des messages entre les instances de consommateur au sein d’un même ID de groupe.

Abonnement en cluster (par défaut)

Chaque message est remis à une seule instance de consommateur du groupe. Tous les consommateurs identifiés par le même ID de groupe consomment un nombre égal de messages.

Par exemple, si un topic contient neuf messages et que le groupe dispose de trois instances de consommateur, chaque instance reçoit trois messages.

// Clustering subscription is the default mode. To set it explicitly:
factoryInfo.setFactoryProperty(ONSFactoryProperty::MessageModel, ONSFactoryProperty::CLUSTERING);

Utilisez l’abonnement en cluster lorsque chaque message ne doit être traité qu’une seule fois, par exemple pour le traitement des commandes ou la répartition des tâches.

Abonnement en diffusion

Chaque message est remis à toutes les instances de consommateur du groupe. Chaque consommateur reçoit une copie intégrale de tous les messages.

Par exemple, si un topic contient neuf messages et que le groupe dispose de trois instances de consommateur, chaque instance reçoit les neuf messages.

// Set broadcasting subscription mode.
factoryInfo.setFactoryProperty(ONSFactoryProperty::MessageModel, ONSFactoryProperty::BROADCASTING);

Privilégiez l’abonnement en diffusion lorsque chaque instance a besoin des mêmes données, par exemple pour l’actualisation du cache local ou la synchronisation de la configuration.

Remarque
  • Toutes les instances de consommateur partageant le même ID de groupe doivent utiliser des abonnements cohérents. Pour plus d’informations, consultez Cohérence de l’abonnement.

  • L’abonnement en diffusion ne prend pas en charge les messages ordonnés, le suivi de la progression de la consommation ni la réinitialisation de l’offset du consommateur. Pour plus d’informations, consultez Consommation en cluster et consommation en diffusion.

Exemple de code

L’exemple suivant crée un PushConsumer qui s’abonne à deux topics, traite les messages entrants via un callback d’écoute et gère le cycle de vie du consommateur.

#include "ONSFactory.h"

#include <iostream>
#include <thread>
#include <mutex>

using namespace ons;

// Mutex to synchronize access to shared resources across listener threads.
std::mutex console_mtx;

class ExampleMessageListener : public MessageListener {
public:
    Action consume(Message& message, ConsumeContext& context) {
        // Return CommitMessage after successful processing.
        // Return ReconsumeLater if processing fails or you want to consume
        // the message again -- the broker redelivers the message after a
        // predefined interval.
        std::lock_guard<std::mutex> lk(console_mtx);
        std::cout << "Received a message. Topic: " << message.getTopic() << ", MsgId: "
        << message.getMsgID() << std::endl;
        return CommitMessage;
    }
};

int main(int argc, char* argv[]) {
    std::cout << "=======Before consuming messages=======" << std::endl;
    ONSFactoryProperty factoryInfo;

    // Group ID created in the ApsaraMQ for RocketMQ console.
    // ApsaraMQ for RocketMQ instances use group IDs instead of producer IDs
    // and consumer IDs. This parameter maintains backward compatibility.
    factoryInfo.setFactoryProperty(ONSFactoryProperty::ConsumerId, "GID_XXX");

    // Retrieve credentials from environment variables.
    // Make sure ALIBABA_CLOUD_ACCESS_KEY_ID and ALIBABA_CLOUD_ACCESS_KEY_SECRET
    // are set before running this program.
    factoryInfo.setFactoryProperty(ONSFactoryProperty::AccessKey, getenv("ALIBABA_CLOUD_ACCESS_KEY_ID"));
    factoryInfo.setFactoryProperty(ONSFactoryProperty::SecretKey, getenv("ALIBABA_CLOUD_ACCESS_KEY_SECRET"));

    // Endpoint of the ApsaraMQ for RocketMQ instance.
    // Get this value from the ApsaraMQ for RocketMQ console.
    factoryInfo.setFactoryProperty(ONSFactoryProperty::NAMESRV_ADDR, "http://xxxxxxxxxxxxxxxx.aliyuncs.com:80");

    PushConsumer *consumer = ONSFactory::getInstance()->createPushConsumer(factoryInfo);

    // Subscribe to messages with a specific tag in topic-1.
    const char* topic_1 = "topic-1";
    const char* tag_1 = "tag-1";

    // Subscribe to all messages in topic-2 by using the wildcard tag "*".
    const char* topic_2 = "topic-2";
    const char* tag_2 = "*";

    // Register the listener and subscribe to both topics.
    ExampleMessageListener * message_listener = new ExampleMessageListener();
    consumer->subscribe(topic_1, tag_1, message_listener);
    consumer->subscribe(topic_2, tag_2, message_listener);

    // Start the consumer. Messages begin arriving after this call.
    consumer->start();

    // Keep the main thread alive. In production, replace this with
    // your application's main loop or signal handler.
    std::this_thread::sleep_for(std::chrono::milliseconds(60 * 1000));
    consumer->shutdown();
    delete message_listener;
    std::cout << "=======After consuming messages======" << std::endl;
    return 0;
}

Étapes suivantes