Cette rubrique fournit des exemples de code permettant d'envoyer et de recevoir des messages transactionnels à l'aide du SDK client TCP pour C++ de la Community Edition.
ApsaraMQ for RocketMQ propose une fonctionnalité de traitement des transactions distribuées similaire aux modes XA (eXtended Architecture) et Open XA. Cette fonctionnalité garantit la cohérence des données dans ApsaraMQ for RocketMQ.
Fonctionnement des messages transactionnels
Le diagramme suivant illustre les interactions entre le producteur, le broker et la transaction locale lors de la livraison des messages transactionnels.

Pour en savoir plus sur le modèle de messagerie transactionnelle, consultez la section Messages transactionnels.
Envoyer des messages transactionnels
-
Copiez le code suivant dans le fichier TransProducerDemo.cpp. Modifiez les paramètres appropriés, exécutez la commande g++ pour compiler le code, puis générez un fichier exécutable.
g++ -o trans_producer_demo -std=c++11 -lz -lrocketmq TransProducerDemo.cpp -
L'exemple de code ci-dessous montre comment envoyer des messages transactionnels à l'aide du SDK client TCP pour C++ de la Community Edition :
#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: LocalTransactionState executeLocalTransaction(const MQMessage &msg, void *arg) { // Execute the local transaction. If the local transaction is executed, COMMIT_MESSAGE is returned. If the local transaction fails to be executed, ROLLBACK_MESSAGE is returned. If the execution status of the local transaction is unknown, UNKNOWN is returned. // If UNKNOWN is returned, the scheduled task to query the status of the local transaction is triggered. std::cout << "Execute Local Transaction,Received Message Topic:" << msg.getTopic() << ", MsgId:" << msg.getBody() << std::endl; return UNKNOWN; } LocalTransactionState checkLocalTransaction(const MQMessageExt &msg) { // Query the execution status of the local transaction. If the local transaction is executed, COMMIT_MESSAGE is returned. If the local transaction fails to be executed, ROLLBACK_MESSAGE is returned. If the execution status of the local transaction is unknown, UNKNOWN is returned. // If UNKNOWN is returned, wait until the next scheduled task to query the status of the local transaction is triggered. std::cout << "Check Local Transaction,Received Message Topic:" << msg.getTopic() << ", MsgId:" << msg.getMsgId() << std::endl; return COMMIT_MESSAGE; } }; int main() { std::cout << "=======Before sending messages=======" << std::endl; // The ID of the group for which you applied in the ApsaraMQ for RocketMQ console. TransactionMQProducer producer("GID_XXXXXXXXXXXXXXXX"); // The TCP endpoint that you obtained from the Instance Details page in the ApsaraMQ for RocketMQ console. producer.setNamesrvAddr("http://MQ_XXXXXXXXXXXX.mq-internet-access.mq-internet.aliyuncs.com:80"); // Make sure that the environment variables ALIBABA_CLOUD_ACCESS_KEY_ID and ALIBABA_CLOUD_ACCESS_KEY_SECRET are configured. // ALIBABA_CLOUD_ACCESS_KEY_ID and ALIBABA_CLOUD_ACCESS_KEY_SECRET specify the AccessKey ID and AccessKey secret of your Alibaba Cloud account, which are used for identity verification. // The user channel. Default value: ALIYUN. producer.setSessionCredentials(getenv("ALIBABA_CLOUD_ACCESS_KEY_ID"), getenv("ALIBABA_CLOUD_ACCESS_KEY_SECRET"), "ALIYUN"); // The transaction listener. ExampleTransactionListener *exampleTransactionListener = new ExampleTransactionListener(); producer.setTransactionListener(exampleTransactionListener); // After you configure the required parameters, start the producer. producer.start(); auto start = std::chrono::system_clock::now(); int count = 3; for (int i = 0; i < count; ++i) { // The topic for which you applied in the ApsaraMQ for RocketMQ console. MQMessage msg("YOUR TRANSACTION TOPIC", "HiTAG", "Hello,CPP SDK, Transaction Message."); try { SendResult sendResult = producer.sendMessageInTransaction(msg, &i); std::cout << "SendResult:" << sendResult.getSendStatus() << ", Message ID: " << sendResult.getMsgId() << std::endl; this_thread::sleep_for(chrono::seconds(1)); } catch (MQException e) { std::cout << "ErrorCode: " << e.GetError() << " Exception:" << e.what() << std::endl; } } auto interval = std::chrono::system_clock::now() - start; std::cout << "Send " << count << " messages OK, costs " << std::chrono::duration_cast<std::chrono::milliseconds>(interval).count() << "ms" << std::endl; std::cout << "Wait for local transaction check..... " << std::endl; for (int i = 0; i < 6; ++i) { this_thread::sleep_for(chrono::seconds(10)); std::cout << "Running "<< i*10 + 10 << " Seconds......"<< std::endl; } producer.shutdown(); std::cout << "=======After sending messages=======" << std::endl; return 0; }
Consommer des messages transactionnels
-
L'exemple de code ci-dessous montre comment consommer des messages transactionnels à l'aide du SDK client TCP pour C++ de la Community Edition :
#include <iostream> #include <thread> #include "DefaultMQPushConsumer.h" using namespace rocketmq; // Implement a message listener to process incoming messages 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; // Group ID from the ApsaraMQ for RocketMQ console DefaultMQPushConsumer *consumer = new DefaultMQPushConsumer("<your-group-id>"); // TCP endpoint from the Instance Details page consumer->setNamesrvAddr("<your-endpoint>"); // Authenticate with AccessKey credentials stored in environment variables. // User channel default: ALIYUN. consumer->setSessionCredentials( getenv("ALIBABA_CLOUD_ACCESS_KEY_ID"), getenv("ALIBABA_CLOUD_ACCESS_KEY_SECRET"), "ALIYUN" ); // Subscribe to the topic and register a listener ExampleMessageListener *messageListener = new ExampleMessageListener(); consumer->subscribe("<your-topic>", "*"); consumer->registerMessageListener(messageListener); // Start the consumer after subscriptions are configured. // All consumers in the same group must have identical subscriptions. consumer->start(); // Keep the main thread alive for 60 seconds to receive messages std::this_thread::sleep_for(std::chrono::milliseconds(60 * 1000)); consumer->shutdown(); std::cout << "=======After consuming messages======" << std::endl; return 0; }