Tous les produits
Search
Centre de documentation

ApsaraMQ for RocketMQ:Send and receive normal messages by using the C++ HTTP SDK

Dernière mise à jour :Aug 09, 2026

Cette rubrique fournit des exemples de code C++ pour l'envoi et la réception de messages normaux à l'aide du SDK client HTTP ApsaraMQ for RocketMQ. Les messages normaux sont des messages standard fournis par ApsaraMQ for RocketMQ : contrairement aux messages planifiés, différés, ordonnés et transactionnels, ils ne comportent aucune sémantique de livraison particulière. Utilisez les messages normaux pour la messagerie générale lorsque l'ordre, le timing et les garanties transactionnelles ne sont pas requis.

Prérequis

Avant de commencer, effectuez les opérations suivantes :

  • Installez le SDK for C++. Pour plus d'informations, consultez la section Préparer l'environnement.

  • Créez dans la console ApsaraMQ for RocketMQ les ressources que vous souhaitez spécifier dans le code. Ces ressources incluent des instances, des topics et des groupes de consommateurs. Pour plus d'informations, consultez la section Créer des ressources.

  • Récupérez la paire AccessKey de votre compte Alibaba Cloud. Pour plus d'informations, consultez la section Créer une AccessKey.

Envoyer des messages normaux

Initialisez un objet MQClient avec votre endpoint HTTP et votre paire AccessKey, puis publiez des messages sur un topic.

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

Espace réservé Description Où le trouver
${HTTP_ENDPOINT} Endpoint d'accès HTTP Section HTTP Endpoint de la page Instance Details
${TOPIC} Nom du topic cible Créé dans la console ApsaraMQ for RocketMQ
${INSTANCE_ID} ID de l'instance. Définissez cette valeur sur une chaîne vide si l'instance n'a pas de namespace Page Instance Details de la console
#include <fstream>
#include <time.h>
#include "mq_http_sdk/mq_client.h"

using namespace std;
using namespace mq::http::sdk;

int main() {

    MQClient mqClient(
            // The HTTP endpoint, found in the HTTP Endpoint section
            // of the Instance Details page in the ApsaraMQ for RocketMQ console.
            "${HTTP_ENDPOINT}",
            // The AccessKey ID and AccessKey secret for authentication.
            // Read from environment variables to avoid hardcoding credentials.
            System.getenv("ALIBABA_CLOUD_ACCESS_KEY_ID"),
            System.getenv("ALIBABA_CLOUD_ACCESS_KEY_SECRET")
            );

    // The topic to publish messages to. Create the topic in the ApsaraMQ for RocketMQ console.
    string topic = "${TOPIC}";
    // The instance ID. If the instance has no namespace, set this to an empty string.
    // Check the namespace on the Instance Details page in the console.
    string instanceId = "${INSTANCE_ID}";

    MQProducerPtr producer;
    if (instanceId == "") {
        producer = mqClient.getProducerRef(topic);
    } else {
        producer = mqClient.getProducerRef(instanceId, topic);
    }

    try {
        // Send four messages in a loop.
        for (int i = 0; i < 4; i++)
        {
            PublishMessageResponse pmResp;

            // The message body.
            TopicMessage pubMsg("Hello, mq!have key!");
            // Set a custom property on the message.
            pubMsg.putProperty("a",std::to_string(i));
            // The message key.
            pubMsg.setMessageKey("MessageKey" + std::to_string(i));
            producer->publishMessage(pubMsg, pmResp);

            cout << "Publish mq message success. Topic is: " << topic
                << ", msgId is:" << pmResp.getMessageId()
                << ", bodyMD5 is:" << pmResp.getMessageBodyMD5() << endl;
        }
    } catch (MQServerException& me) {
        cout << "Request Failed: " + me.GetErrorCode() << ", requestId is:" << me.GetRequestId() << endl;
        return -1;
    } catch (MQExceptionBase& mb) {
        cout << "Request Failed: " + mb.ToString() << endl;
        return -2;
    }

    return 0;
}

Points clés :

  • Définissez les variables d'environnement ALIBABA_CLOUD_ACCESS_KEY_ID et ALIBABA_CLOUD_ACCESS_KEY_SECRET avant d'exécuter le code.

  • Chaque message prend en charge des propriétés personnalisées (putProperty) et une clé de message (setMessageKey).

  • Si l'instance n'a pas de namespace, transmettez une chaîne vide pour instanceId. Le SDK sélectionne automatiquement la bonne surcharge getProducerRef.

S'abonner à des messages normaux

Fonctionnement de la consommation

Le consommateur utilise le sondage long (long polling) HTTP pour récupérer les messages :

  1. Le consommateur envoie une requête au broker.

  2. Si des messages sont disponibles, le broker répond immédiatement.

  3. Si aucun message n'est disponible, le broker maintient la requête ouverte pendant la durée de sondage spécifiée (maximum : 30 secondes) et répond dès qu'un message arrive.

  4. Après avoir traité chaque lot, le consommateur envoie un accusé de réception (ACK) au broker.

  5. Si le broker ne reçoit pas d'ACK avant l'échéance NextConsumeTime, il redistribue le message.

Paramètres du consommateur

Paramètre Description Valeur d'exemple Maximum
Taille du lot Nombre maximal de messages consommés par requête 3 16
Durée de sondage Temps d'attente du broker pour les messages (en secondes) 3 30

Exemple de code

Remplacez ces espaces réservés en plus de ceux indiqués dans la section Envoyer des messages normaux :

Espace réservé Description Où le trouver
${GROUP_ID} ID du groupe de consommateurs Créé dans la console ApsaraMQ for RocketMQ
#include <vector>
#include <fstream>
#include "mq_http_sdk/mq_client.h"

#ifdef _WIN32
#include <windows.h>
#else
#include <unistd.h>
#endif

using namespace std;
using namespace mq::http::sdk;

int main() {

    MQClient mqClient(
            // The HTTP endpoint, found in the HTTP Endpoint section
            // of the Instance Details page in the ApsaraMQ for RocketMQ console.
            "${HTTP_ENDPOINT}",
            // The AccessKey ID and AccessKey secret for authentication.
            // Read from environment variables to avoid hardcoding credentials.
            System.getenv("ALIBABA_CLOUD_ACCESS_KEY_ID"),
            System.getenv("ALIBABA_CLOUD_ACCESS_KEY_SECRET")
            );

    // The topic to consume messages from. Create the topic in the ApsaraMQ for RocketMQ console.
    string topic = "${TOPIC}";
    // The consumer group ID. Create the consumer group in the ApsaraMQ for RocketMQ console.
    string groupId = "${GROUP_ID}";
    // The instance ID. If the instance has no namespace, set this to an empty string.
    // Check the namespace on the Instance Details page in the console.
    string instanceId = "${INSTANCE_ID}";

    MQConsumerPtr consumer;
    if (instanceId == "") {
        consumer = mqClient.getConsumerRef(topic, groupId);
    } else {
        consumer = mqClient.getConsumerRef(instanceId, topic, groupId, "");
    }

    do {
        try {
            std::vector<Message> messages;
            // Consume messages in long polling mode.
            // If no messages are available, the broker holds the request for the specified
            // polling duration and responds immediately when a message arrives.
            consumer->consumeMessage(
                    3, // Batch size: max messages per request. Maximum allowed: 16.
                    3, // Polling duration in seconds. Maximum allowed: 30.
                    messages
            );
            cout << "Consume: " << messages.size() << " Messages!" << endl;

            // Process each message and collect receipt handles for acknowledgment.
            std::vector<std::string> receiptHandles;
            for (std::vector<Message>::iterator iter = messages.begin();
                    iter != messages.end(); ++iter)
            {
                cout << "MessageId: " << iter->getMessageId()
                    << " PublishTime: " << iter->getPublishTime()
                    << " Tag: " << iter->getMessageTag()
                    << " Body: " << iter->getMessageBody()
                    << " FirstConsumeTime: " << iter->getFirstConsumeTime()
                    << " NextConsumeTime: " << iter->getNextConsumeTime()
                    << " ConsumedTimes: " << iter->getConsumedTimes()
                    << " Properties: " << iter->getPropertiesAsString()
                    << " Key: " << iter->getMessageKey() << endl;
                receiptHandles.push_back(iter->getReceiptHandle());
            }

            // Send an ACK to the broker. If the broker does not receive an ACK before
            // NextConsumeTime, it redelivers the message. Each receipt handle carries
            // a unique timestamp that expires after the NextConsumeTime window.
            AckMessageResponse bdmResp;
            consumer->ackMessage(receiptHandles, bdmResp);
            if (!bdmResp.isSuccess()) {
                // Log failed ACK items. A handle that has expired causes ACK failure,
                // and the broker redelivers the corresponding message.
                const std::vector<AckMessageFailedItem>& failedItems =
                    bdmResp.getAckMessageFailedItem();
                for (std::vector<AckMessageFailedItem>::const_iterator iter = failedItems.begin();
                        iter != failedItems.end(); ++iter)
                {
                    cout << "AckFailedItem: " << iter->errorCode
                        << "  " << iter->receiptHandle << endl;
                }
            } else {
                cout << "Ack: " << messages.size() << " messages suc!" << endl;
            }
        } catch (MQServerException& me) {
            if (me.GetErrorCode() == "MessageNotExist") {
                cout << "No message to consume! RequestId: " + me.GetRequestId() << endl;
                continue;
            }
            cout << "Request Failed: " + me.GetErrorCode() + ".RequestId: " + me.GetRequestId() << endl;
#ifdef _WIN32
            Sleep(2000);
#else
            usleep(2000 * 1000);
#endif
        } catch (MQExceptionBase& mb) {
            cout << "Request Failed: " + mb.ToString() << endl;
#ifdef _WIN32
            Sleep(2000);
#else
            usleep(2000 * 1000);
#endif
        }

    } while(true);
}

Points clés :

  • Le consommateur s'exécute dans une boucle infinie, effectuant un sondage continu. Lorsqu'aucun message n'est disponible, il intercepte l'exception MessageNotExist et réessaie.

  • Après avoir traité un lot, appelez ackMessage pour confirmer la réception. Les messages non acquittés sont redistribués après l'échéance NextConsumeTime.

  • Chaque handle de réception comporte un horodatage unique. Un handle expiré entraîne un échec de l'ACK, ce qui déclenche une redistribution.

  • En cas d'erreurs autres que MessageNotExist, le consommateur marque une pause de 2 secondes avant de réessayer afin d'éviter les boucles serrées.