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_IDetALIBABA_CLOUD_ACCESS_KEY_SECRETavant 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 surchargegetProducerRef.
S'abonner à des messages normaux
Fonctionnement de la consommation
Le consommateur utilise le sondage long (long polling) HTTP pour récupérer les messages :
Le consommateur envoie une requête au broker.
Si des messages sont disponibles, le broker répond immédiatement.
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.
Après avoir traité chaque lot, le consommateur envoie un accusé de réception (ACK) au broker.
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
MessageNotExistet réessaie.Après avoir traité un lot, appelez
ackMessagepour confirmer la réception. Les messages non acquittés sont redistribués après l'échéanceNextConsumeTime.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.