Les messages normaux sont le type de message le plus basique dans ApsaraMQ for RocketMQ. Contrairement aux messages planifiés, différés, ordonnés et transactionnels, les messages normaux n'ont aucune sémantique de livraison particulière. Utilisez des messages normaux lorsque votre application nécessite une livraison fiable sans contraintes d'ordre ni de délai.
Les exemples suivants montrent comment envoyer et consommer des messages normaux avec le SDK client TCP pour C++ (édition communautaire).
Prérequis
Avant de commencer, assurez-vous d'avoir :
Installé la bibliothèque dynamique C++ pour le SDK RocketMQ. Pour plus d'informations, consultez Installation de la bibliothèque dynamique C++.
Créé une paire AccessKey pour votre compte Alibaba Cloud. Pour plus d'informations, consultez Création d'une paire AccessKey.
Défini les variables d'environnement
ALIBABA_CLOUD_ACCESS_KEY_IDetALIBABA_CLOUD_ACCESS_KEY_SECRETavec votre ID AccessKey et votre secret AccessKey.
Envoi de messages normaux
Étape 1 : Créer le fichier du producteur
Enregistrez le code suivant sous le nom ProducerDemo.cpp. Remplacez les espaces réservés par vos valeurs réelles :
| Espace réservé | Description | Exemple |
|---|---|---|
<your-group-id> |
L'ID de groupe créé dans la console ApsaraMQ for RocketMQ | GID_demo_producer |
<your-endpoint> |
L'endpoint TCP affiché sur la page Détails de l'instance dans la console ApsaraMQ for RocketMQ | http://MQ_INST_xxxxx.mq-internet-access.mq-internet.aliyuncs.com:80 |
<your-topic> |
Le topic créé dans la console ApsaraMQ for RocketMQ | demo_topic |
#include <iostream>
#include <chrono>
#include <thread>
#include "DefaultMQProducer.h"
using namespace std;
using namespace rocketmq;
int main() {
std::cout << "=======Before sending messages=======" << std::endl;
// Group ID from the ApsaraMQ for RocketMQ console
DefaultMQProducer producer("<your-group-id>");
// TCP endpoint from the Instance Details page
producer.setNamesrvAddr("<your-endpoint>");
// Authenticate with AccessKey credentials stored in environment variables.
// User channel default: ALIYUN.
producer.setSessionCredentials(
getenv("ALIBABA_CLOUD_ACCESS_KEY_ID"),
getenv("ALIBABA_CLOUD_ACCESS_KEY_SECRET"),
"ALIYUN"
);
// Start the producer
producer.start();
auto start = std::chrono::system_clock::now();
int count = 32;
for (int i = 0; i < count; ++i) {
// Construct a message with topic, tag, and body
MQMessage msg("<your-topic>", "HiTAG", "HelloCPPSDK.");
try {
SendResult sendResult = producer.send(msg);
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;
producer.shutdown();
std::cout << "=======After sending messages=======" << std::endl;
return 0;
}
Étape 2 : Compiler et exécuter le producteur
g++ -o producer_demo -std=c++11 -lz -lrocketmq ProducerDemo.cpp
./producer_demo
Sortie attendue :
SendResult:0, Message ID: C0A8XXXXXXXXXXXXXXXXXXXX
Un statut d'envoi égal à 0 indique que l'opération a réussi. Chaque message affiche son ID unique.
Consommation de messages normaux
Étape 1 : Créer le fichier du consommateur
Enregistrez le code suivant sous le nom ConsumerDemo.cpp. Utilisez les mêmes valeurs d'espace réservé que dans l'exemple du producteur ci-dessus.
#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;
}
Étape 2 : Compiler et exécuter le consommateur
g++ -o consumer_demo -std=c++11 -lz -lrocketmq ConsumerDemo.cpp
./consumer_demo
Sortie attendue :
Received Message Topic:demo_topic, MsgId:C0A8XXXXXXXXXXXXXXXXXXXX
Bonnes pratiques
Cohérence de l'ID de groupe : Tous les consommateurs appartenant au même groupe doivent s'abonner aux mêmes topics avec des filtres de tags identiques. Des abonnements incompatibles entraînent un comportement de livraison indéfini.
Cycle de vie du consommateur : Cet exemple s'exécute pendant 60 secondes avant de s'arrêter. En production, maintenez le processus du consommateur actif aussi longtemps que votre application doit recevoir des messages.
Gestion des erreurs : Le producteur intercepte l'exception
MQExceptionet affiche le code d'erreur. Adaptez ce modèle en ajoutant des mécanismes de nouvelle tentative ou d'alerte selon vos besoins.