ApsaraMQ for RocketMQ prend en charge deux types de messages à différé temporel via le SDK client HTTP pour C++ :
Messages différés : livrés après un délai spécifié à compter de l'envoi. Par exemple, avec un délai de 10 secondes, le consommateur reçoit le message 10 secondes après son envoi.
Messages planifiés : livrés à un instant précis. Par exemple, un message programmé pour 14 h 00 est livré à 14 h 00.
Ces deux types utilisent la même API. Définissez StartDeliverTime sur un horodatage Unix en millisecondes correspondant à l'instant de livraison souhaité :
Message différé : heure actuelle + durée du délai. Exemple :
time(NULL) * 1000 + 10 * 1000pour une livraison après 10 secondes.Message planifié : l'instant cible de livraison converti en horodatage Unix (millisecondes).
Pour en savoir plus sur ces types de messages, consultez la rubrique Messages planifiés et messages différés.
Prérequis
Avant de commencer, assurez-vous d'avoir :
Le SDK client HTTP C++ installé. Pour plus d'informations, consultez la rubrique Préparer l'environnement.
Une instance ApsaraMQ for RocketMQ, ainsi qu'un topic et un groupe de consommateurs créés dans la console ApsaraMQ for RocketMQ.
Une paire AccessKey pour votre compte Alibaba Cloud. Pour plus d'informations, consultez la rubrique Créer une paire AccessKey.
Envoyer des messages planifiés et différés
L'extrait de code suivant envoie quatre messages différés, chacun livré 10 secondes après l'envoi. Pour envoyer un message planifié à la place, définissez StartDeliverTime sur l'instant de livraison cible sous forme d'horodatage Unix en millisecondes.
Remplacez les espaces réservés suivants par vos valeurs réelles :
| Espace réservé | Description | Emplacement |
|---|---|---|
${HTTP_ENDPOINT} |
Endpoint HTTP de l'instance | Page Instance Details > section HTTP Endpoint dans la console ApsaraMQ for RocketMQ |
${TOPIC} |
Nom du topic | Console ApsaraMQ for RocketMQ |
${INSTANCE_ID} |
ID de l'instance. Laissez vide si l'instance n'utilise pas de namespace | Page Instance Details dans la console ApsaraMQ for RocketMQ |
#include <fstream>
#include <time.h>
#include "mq_http_sdk/mq_client.h"
using namespace std;
using namespace mq::http::sdk;
int main() {
MQClient mqClient(
// HTTP endpoint of the instance.
"${HTTP_ENDPOINT}",
// Make sure that the environment variables ALIBABA_CLOUD_ACCESS_KEY_ID
// and ALIBABA_CLOUD_ACCESS_KEY_SECRET are configured.
// AccessKey ID for authentication.
System.getenv("ALIBABA_CLOUD_ACCESS_KEY_ID"),
// AccessKey secret for authentication.
System.getenv("ALIBABA_CLOUD_ACCESS_KEY_SECRET")
);
// Topic to send messages to. Create this topic in the ApsaraMQ for RocketMQ console.
string topic = "${TOPIC}";
// Instance ID. If the instance has a namespace, specify the instance ID.
// If not, set this to an empty string.
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;
// Message body.
TopicMessage pubMsg("Hello, mq!have key!");
// Custom message property.
pubMsg.putProperty("a",std::to_string(i));
// Message key.
pubMsg.setMessageKey("MessageKey" + std::to_string(i));
// Deliver after a 10-second delay.
// StartDeliverTime is a millisecond-level Unix timestamp.
// For a scheduled message, set this to the target delivery time
// as a millisecond-level Unix timestamp.
pubMsg.setStartDeliverTime(time(NULL) * 1000 + 10 * 1000);
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;
}
S'abonner aux messages planifiés et différés
Les messages planifiés et différés sont consommés de la même manière que les messages classiques. Le broker conserve chaque message jusqu'à son heure de livraison, puis le rend disponible pour consommation.
L'extrait de code suivant utilise le sondage long (long polling) pour consommer les messages. Dans ce mode, si aucun message n'est disponible, la requête reste ouverte sur le broker pendant la durée spécifiée. Le broker répond immédiatement dès qu'un message arrive.
Remplacez les espaces réservés suivants par vos valeurs réelles :
| Espace réservé | Description | Emplacement |
|---|---|---|
${HTTP_ENDPOINT} |
Endpoint HTTP de l'instance | Page Instance Details > section HTTP Endpoint dans la console ApsaraMQ for RocketMQ |
${TOPIC} |
Nom du topic | Console ApsaraMQ for RocketMQ |
${GROUP_ID} |
ID du groupe de consommateurs | Console ApsaraMQ for RocketMQ |
${INSTANCE_ID} |
ID de l'instance. Laissez vide si l'instance n'utilise pas de namespace | Page Instance Details 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(
// HTTP endpoint of the instance.
"${HTTP_ENDPOINT}",
// Make sure that the environment variables ALIBABA_CLOUD_ACCESS_KEY_ID
// and ALIBABA_CLOUD_ACCESS_KEY_SECRET are configured.
// AccessKey ID for authentication.
System.getenv("ALIBABA_CLOUD_ACCESS_KEY_ID"),
// AccessKey secret for authentication.
System.getenv("ALIBABA_CLOUD_ACCESS_KEY_SECRET")
);
// Topic to consume messages from. Create this topic in the ApsaraMQ for RocketMQ console.
string topic = "${TOPIC}";
// Consumer group ID. Create this in the ApsaraMQ for RocketMQ console.
string groupId = "${GROUP_ID}";
// Instance ID. If the instance has a namespace, specify the instance ID.
// If not, set this to an empty string.
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 message is available, the request is held on the broker
// for the specified duration until a message arrives.
consumer->consumeMessage(
3, // Maximum messages per request (max: 16).
3, // Long polling timeout in seconds (max: 30).
messages
);
cout << "Consume: " << messages.size() << " Messages!" << endl;
// Process consumed messages.
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());
}
// ACK the consumed messages.
// If the broker does not receive an ACK before NextConsumeTime,
// it redelivers the message. Each consumption attempt generates
// a new receipt handle with a unique timestamp.
AckMessageResponse bdmResp;
consumer->ackMessage(receiptHandles, bdmResp);
if (!bdmResp.isSuccess()) {
// Log failed ACKs. A timed-out receipt handle prevents
// the broker from receiving the ACK.
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);
}
Voir aussi
Messages planifiés et messages différés – Concepts et limites relatifs aux messages planifiés et différés
Préparer l'environnement – Configurer le SDK client HTTP C++