ApsaraMQ for RocketMQ 5.x est compatible avec les clients SDK RocketMQ 3.x et 4.x. Cette rubrique fournit des exemples de code en C++ pour l'envoi et la réception de messages standard, ordonnés, planifiés/différés et transactionnels.
Nous vous recommandons d'utiliser les dernières versions des SDK RocketMQ 5.x. Ces SDK sont entièrement compatibles avec les brokers ApsaraMQ for RocketMQ 5.x et offrent davantage de fonctionnalités ainsi que des améliorations. Pour plus d'informations, consultez la page Informations sur la version du SDK C++.
Alibaba Cloud ne maintient que les SDK client RocketMQ 4.x, 3.x et TCP. Nous vous conseillons de les utiliser uniquement pour les activités existantes.
Prérequis
Avant de commencer, assurez-vous d'avoir :
Installé la bibliothèque dynamique C++. Pour plus de détails, consultez la page Installer la bibliothèque dynamique C++
Créé un topic et un ID de groupe dans la console ApsaraMQ for RocketMQ
Obtenu l'endpoint de l'instance depuis la console ApsaraMQ for RocketMQ
Configuration commune
Tous les exemples de cette rubrique partagent la même configuration de connexion. Remplacez les espaces réservés suivants par vos propres valeurs avant d'exécuter le code.
| Espace réservé | Description | Exemple |
|---|---|---|
<your-group-id> |
ID de groupe créé dans la console | GID_example |
<your-access-point> |
Endpoint de l'instance au format host:port. N'incluez pas http:// ni https://, et n'utilisez pas d'adresse IP résolue. |
rmq-cn-xxx.rmq.aliyuncs.com:8080 |
<your-topic> |
Topic créé dans la console | normal_topic_01 |
<instance-username> |
Nom d'utilisateur de l'instance, disponible dans l'onglet Intelligent Authentication de la page Access Control | N/A |
<instance-password> |
Mot de passe de l'instance, disponible dans l'onglet Intelligent Authentication de la page Access Control | N/A |
Authentification :
Accès Internet : nom d'utilisateur et mot de passe requis.
Accès VPC : aucun nom d'utilisateur ni mot de passe nécessaire.
Instance Serverless via Internet : nom d'utilisateur et mot de passe requis.
Instance Serverless dans un VPC avec accès sans authentification activé : aucun nom d'utilisateur ni mot de passe nécessaire.
Ne spécifiez pas l'ID de l'instance lorsque vous utilisez le SDK RocketMQ 3.x ou 4.x pour C++ avec une instance ApsaraMQ for RocketMQ 5.0. Cela entraînerait des échecs de connexion.
Configuration commune du producteur :
#include <iostream>
#include <chrono>
#include <thread>
#include "DefaultMQProducer.h"
using namespace std;
using namespace rocketmq;
// Initialize the producer with your group ID
DefaultMQProducer producer("<your-group-id>");
// Set the endpoint (domain:port only, no http/https prefix)
producer.setNamesrvAddr("<your-access-point>");
// Set credentials (required for Internet access; skip for VPC access)
producer.setSessionCredentials("<instance-username>", "<instance-password>", "ALIYUN");
producer.start();
Configuration commune du consommateur :
#include <iostream>
#include <thread>
#include "DefaultMQPushConsumer.h"
using namespace rocketmq;
// Initialize the push consumer with your group ID
DefaultMQPushConsumer *consumer = new DefaultMQPushConsumer("<your-group-id>");
// Set the endpoint (domain:port only, no http/https prefix)
consumer->setNamesrvAddr("<your-access-point>");
// Set credentials (required for Internet access; skip for VPC access)
consumer->setSessionCredentials("<instance-username>", "<instance-password>", "ALIYUN");
Avant de démarrer un consommateur :
Configurez tous les abonnements avant d'appeler
start().Assurez-vous que tous les consommateurs d'un même groupe utilisent des abonnements identiques.
Messages standard
Utilisez les messages standard pour la messagerie générale qui ne nécessite ni ordre spécifique, ni livraison différée, ni prise en charge des transactions.
Envoyer des messages standard
#include <iostream>
#include <chrono>
#include <thread>
#include "DefaultMQProducer.h"
using namespace std;
using namespace rocketmq;
int main() {
cout << "=======Before sending messages=======" << endl;
DefaultMQProducer producer("<your-group-id>");
producer.setNamesrvAddr("<your-access-point>");
producer.setSessionCredentials("<instance-username>", "<instance-password>", "ALIYUN");
producer.start();
auto start = chrono::system_clock::now();
int count = 32;
for (int i = 0; i < count; ++i) {
MQMessage msg("<your-topic>", "HiTAG", "HelloCPPSDK.");
try {
SendResult sendResult = producer.send(msg);
cout << "SendResult:" << sendResult.getSendStatus()
<< ", Message ID: " << sendResult.getMsgId() << endl;
this_thread::sleep_for(chrono::seconds(1));
} catch (MQException& e) {
cout << "ErrorCode: " << e.GetError()
<< " Exception:" << e.what() << endl;
}
}
auto interval = chrono::system_clock::now() - start;
cout << "Send " << count << " messages OK, costs "
<< chrono::duration_cast<chrono::milliseconds>(interval).count()
<< "ms" << endl;
producer.shutdown();
cout << "=======After sending messages=======" << endl;
return 0;
}
S'abonner aux messages standard
#include <iostream>
#include <thread>
#include "DefaultMQPushConsumer.h"
using namespace rocketmq;
class ExampleMessageListener : public MessageListenerConcurrently {
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;
DefaultMQPushConsumer *consumer = new DefaultMQPushConsumer("<your-group-id>");
consumer->setNamesrvAddr("<your-access-point>");
consumer->setSessionCredentials("<instance-username>", "<instance-password>", "ALIYUN");
// Register the listener and subscribe to the topic
ExampleMessageListener *messageListener = new ExampleMessageListener();
consumer->subscribe("<your-topic>", "*");
consumer->registerMessageListener(messageListener);
// Start after subscriptions are configured
consumer->start();
// Keep the main thread alive
std::this_thread::sleep_for(std::chrono::milliseconds(60 * 1000));
consumer->shutdown();
std::cout << "=======After consuming messages======" << std::endl;
return 0;
}
Messages ordonnés
Utilisez les messages ordonnés lorsque les messages appartenant à une même entité métier doivent être traités séquentiellement, comme c'est le cas pour les mises à jour de statut de commande ou les modifications de stock.
Les messages ordonnés acheminent les messages partageant la même clé de partition vers la même file d'attente via un MessageQueueSelector, ce qui garantit l'ordre de traitement.
Envoyer des messages ordonnés
#include <iostream>
#include <chrono>
#include <thread>
#include "DefaultMQProducer.h"
using namespace std;
using namespace rocketmq;
class ExampleSelectMessageQueueByHash : public MessageQueueSelector {
public:
MQMessageQueue select(const std::vector<MQMessageQueue> &mqs,
const MQMessage &msg, void *arg) {
// Route messages to a queue based on the partition key
int orderId = *static_cast<int *>(arg);
int index = orderId % mqs.size();
return mqs[index];
}
};
int main() {
cout << "=======Before sending messages=======" << endl;
DefaultMQProducer producer("<your-group-id>");
producer.setNamesrvAddr("<your-access-point>");
producer.setSessionCredentials("<instance-username>", "<instance-password>", "ALIYUN");
producer.start();
auto start = chrono::system_clock::now();
int count = 32;
ExampleSelectMessageQueueByHash *pSelector = new ExampleSelectMessageQueueByHash();
for (int i = 0; i < count; ++i) {
MQMessage msg("<your-topic>", "HiTAG", "Hello,CPP SDK, Orderly Message.");
try {
// Pass the partition key (i) and the queue selector
SendResult sendResult = producer.send(msg, pSelector, &i, 1, false);
cout << "SendResult:" << sendResult.getSendStatus()
<< ", Message ID: " << sendResult.getMsgId()
<< " MessageQueue:" << sendResult.getMessageQueue().toString()
<< endl;
this_thread::sleep_for(chrono::seconds(1));
} catch (MQException& e) {
cout << "ErrorCode: " << e.GetError()
<< " Exception:" << e.what() << endl;
}
}
auto interval = chrono::system_clock::now() - start;
cout << "Send " << count << " messages OK, costs "
<< chrono::duration_cast<chrono::milliseconds>(interval).count()
<< "ms" << endl;
producer.shutdown();
cout << "=======After sending messages=======" << endl;
return 0;
}
S'abonner aux messages ordonnés
Le consommateur utilise MessageListenerOrderly au lieu de MessageListenerConcurrently afin de traiter les messages dans l'ordre au sein de chaque file d'attente.
#include <iostream>
#include <thread>
#include "DefaultMQPushConsumer.h"
using namespace rocketmq;
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;
DefaultMQPushConsumer *consumer = new DefaultMQPushConsumer("<your-group-id>");
consumer->setNamesrvAddr("<your-access-point>");
consumer->setSessionCredentials("<instance-username>", "<instance-password>", "ALIYUN");
// Register the orderly listener and subscribe to the topic
ExampleOrderlyMessageListener *messageListener = new ExampleOrderlyMessageListener();
consumer->subscribe("<your-topic>", "*");
consumer->registerMessageListener(messageListener);
// Start after subscriptions are configured
consumer->start();
// Keep the main thread alive
std::this_thread::sleep_for(std::chrono::seconds(60));
consumer->shutdown();
std::cout << "=======After consuming messages======" << std::endl;
return 0;
}
Messages planifiés et différés
Utilisez les messages planifiés ou différés pour livrer des messages à un moment précis ou après un délai, par exemple pour des notifications programmées, des mécanismes de nouvelle tentative ou l'exécution différée de tâches.
Définissez la propriété de message __STARTDELIVERTIME avec un horodatage Unix en millisecondes. Le broker conserve le message jusqu'à l'heure spécifiée.
Si l'horodatage est dans le futur, le broker livre le message à ce moment-là.
Si l'horodatage est dans le passé, le broker livre le message immédiatement.
Envoyer des messages planifiés ou différés
#include <iostream>
#include <chrono>
#include <thread>
#include "DefaultMQProducer.h"
using namespace std;
using namespace rocketmq;
int main() {
cout << "=======Before sending messages=======" << endl;
DefaultMQProducer producer("<your-group-id>");
producer.setNamesrvAddr("<your-access-point>");
producer.setSessionCredentials("<instance-username>", "<instance-password>", "ALIYUN");
producer.start();
auto start = chrono::system_clock::now();
int count = 32;
for (int i = 0; i < count; ++i) {
MQMessage msg("<your-topic>", "HiTAG", "Hello,CPP SDK, Delay Message.");
// Calculate a delivery time 10 seconds from now
chrono::system_clock::duration d = chrono::system_clock::now().time_since_epoch();
chrono::milliseconds mil = chrono::duration_cast<chrono::milliseconds>(d);
long exp = mil.count() + 10000; // Delay: 10,000 ms (10 seconds)
msg.setProperty("__STARTDELIVERTIME", to_string(exp));
cout << "Now: " << mil.count() << " Exp:" << exp << endl;
try {
SendResult sendResult = producer.send(msg);
cout << "SendResult:" << sendResult.getSendStatus()
<< ", Message ID: " << sendResult.getMsgId() << endl;
this_thread::sleep_for(chrono::seconds(1));
} catch (MQException& e) {
cout << "ErrorCode: " << e.GetError()
<< " Exception:" << e.what() << endl;
}
}
auto interval = chrono::system_clock::now() - start;
cout << "Send " << count << " messages OK, costs "
<< chrono::duration_cast<chrono::milliseconds>(interval).count()
<< "ms" << endl;
producer.shutdown();
cout << "=======After sending messages=======" << endl;
return 0;
}
S'abonner aux messages planifiés ou différés
Le consommateur utilise MessageListenerConcurrently, le même type d'écouteur que pour les messages standard. Le broker gère la temporisation ; le code du consommateur reste identique.
#include <iostream>
#include <thread>
#include <chrono>
#include "DefaultMQPushConsumer.h"
using namespace rocketmq;
using namespace std;
class ExampleDelayMessageListener : public MessageListenerConcurrently {
public:
ConsumeStatus consumeMessage(const std::vector<MQMessageExt> &msgs) {
for (auto item = msgs.begin(); item != msgs.end(); item++) {
chrono::system_clock::duration d = chrono::system_clock::now().time_since_epoch();
chrono::milliseconds mil = chrono::duration_cast<chrono::milliseconds>(d);
cout << "Now: " << mil.count()
<< " Received Message Topic:" << item->getTopic()
<< ", MsgId:" << item->getMsgId()
<< " DelayTime:" << item->getProperty("__STARTDELIVERTIME")
<< endl;
}
return CONSUME_SUCCESS;
}
};
int main(int argc, char *argv[]) {
cout << "=======Before consuming messages=======" << endl;
DefaultMQPushConsumer *consumer = new DefaultMQPushConsumer("<your-group-id>");
consumer->setNamesrvAddr("<your-access-point>");
consumer->setSessionCredentials("<instance-username>", "<instance-password>", "ALIYUN");
// Register the listener and subscribe to the topic
ExampleDelayMessageListener *messageListener = new ExampleDelayMessageListener();
consumer->subscribe("<your-topic>", "*");
consumer->registerMessageListener(messageListener);
// Start after subscriptions are configured
consumer->start();
// Keep the main thread alive (10 minutes to allow delayed messages to arrive)
this_thread::sleep_for(chrono::seconds(600));
consumer->shutdown();
cout << "=======After consuming messages======" << endl;
return 0;
}
Messages transactionnels
Utilisez les messages transactionnels lorsque la logique métier locale et la livraison des messages doivent réussir ou échouer ensemble, comme lors de transferts de fonds entre comptes.
La messagerie transactionnelle fonctionne en trois étapes :
Le producteur envoie un message semi-validé (half message) au broker.
Le broker invoque
executeLocalTransactioncôté producteur. RenvoyezCOMMIT_MESSAGEpour livrer le message,ROLLBACK_MESSAGEpour l'ignorer, ouUNKNOWNpour différer la décision.Si le résultat est
UNKNOWN, le broker appelle périodiquementcheckLocalTransactionjusqu'à obtenir un résultat définitif.
Envoyer des messages transactionnels
#include <iostream>
#include <chrono>
#include <thread>
#include "TransactionMQProducer.h"
#include "MQClientException.h"
#include "TransactionListener.h"
using namespace std;
using namespace rocketmq;
class ExampleTransactionListener : public TransactionListener {
public:
// Called after the half message is sent
LocalTransactionState executeLocalTransaction(const MQMessage &msg, void *arg) {
cout << "Execute Local Transaction, Received Message Topic:" << msg.getTopic()
<< ", Body:" << msg.getBody() << endl;
// Run local business logic here. Return:
// COMMIT_MESSAGE - deliver the message to consumers
// ROLLBACK_MESSAGE - discard the message
// UNKNOWN - defer; the broker will call checkLocalTransaction later
return UNKNOWN;
}
// Called by the broker when executeLocalTransaction returned UNKNOWN
LocalTransactionState checkLocalTransaction(const MQMessageExt &msg) {
cout << "Check Local Transaction, Received Message Topic:" << msg.getTopic()
<< ", MsgId:" << msg.getMsgId() << endl;
// Query local transaction status and return the result
return COMMIT_MESSAGE;
}
};
int main() {
cout << "=======Before sending messages=======" << endl;
// Use TransactionMQProducer instead of DefaultMQProducer
TransactionMQProducer producer("<your-group-id>");
producer.setNamesrvAddr("<your-access-point>");
producer.setSessionCredentials("<instance-username>", "<instance-password>", "ALIYUN");
// Register the transaction listener before starting the producer
ExampleTransactionListener *exampleTransactionListener = new ExampleTransactionListener();
producer.setTransactionListener(exampleTransactionListener);
producer.start();
auto start = chrono::system_clock::now();
int count = 3;
for (int i = 0; i < count; ++i) {
MQMessage msg("<your-topic>", "HiTAG", "Hello,CPP SDK, Transaction Message.");
try {
SendResult sendResult = producer.sendMessageInTransaction(msg, &i);
cout << "SendResult:" << sendResult.getSendStatus()
<< ", Message ID: " << sendResult.getMsgId() << endl;
this_thread::sleep_for(chrono::seconds(1));
} catch (MQException& e) {
cout << "ErrorCode: " << e.GetError()
<< " Exception:" << e.what() << endl;
}
}
auto interval = chrono::system_clock::now() - start;
cout << "Send " << count << " messages OK, costs "
<< chrono::duration_cast<chrono::milliseconds>(interval).count()
<< "ms" << endl;
// Wait 60 seconds for the broker to call checkLocalTransaction
cout << "Wait for local transaction check..... " << endl;
for (int i = 0; i < 6; ++i) {
this_thread::sleep_for(chrono::seconds(10));
cout << "Running " << i * 10 + 10 << " Seconds......" << endl;
}
producer.shutdown();
cout << "=======After sending messages=======" << endl;
return 0;
}
S'abonner aux messages transactionnels
Abonnez-vous aux messages transactionnels de la même manière qu'aux messages standard. Utilisez MessageListenerConcurrently et DefaultMQPushConsumer. Pour le code complet, consultez la section S'abonner aux messages standard.