Les messages normaux constituent le type de message de base dans ApsaraMQ for RocketMQ. Contrairement aux messages planifiés, différés, ordonnés et transactionnels, les messages normaux n'ont pas de sémantique de livraison particulière.
Les exemples présentés sur cette page utilisent le SDK client HTTP pour Node.js (@aliyunmq/mq-http-sdk) afin de publier et de consommer des messages normaux.
Avant de commencer
Installez le SDK : configurez l'environnement de développement Node.js et installez le package
@aliyunmq/mq-http-sdk. Pour plus d'informations, consultez la rubrique Préparer l'environnement.Créez des ressources : créez une instance, un topic et un groupe de consommateurs dans la console ApsaraMQ for RocketMQ. Pour plus d'informations, consultez la rubrique Créer des ressources.
-
Configurez les identifiants : obtenez une paire AccessKey pour votre compte Alibaba Cloud et exportez les valeurs en tant que variables d'environnement. Pour plus d'informations, consultez la rubrique Créer une paire AccessKey.
export ALIBABA_CLOUD_ACCESS_KEY_ID=<your-access-key-id> export ALIBABA_CLOUD_ACCESS_KEY_SECRET=<your-access-key-secret>
Espaces réservés
Remplacez les espaces réservés suivants dans l'exemple de code par vos valeurs réelles :
| Espace réservé | Description | Où le trouver |
|---|---|---|
<http-endpoint> |
Endpoint d'accès HTTP de votre instance | Page Instance Details > HTTP Endpoint dans la console ApsaraMQ for RocketMQ |
<topic> |
Nom du topic | Page Topics dans la console ApsaraMQ for RocketMQ |
<instance-id> |
ID de l'instance. Définissez cette valeur sur null ou "" si l'instance ne possède pas de namespace. |
Page Instance Details dans la console ApsaraMQ for RocketMQ |
<group-id> |
ID du groupe de consommateurs (consommateur uniquement) | Page Groups dans la console ApsaraMQ for RocketMQ |
Envoyer des messages normaux
-
Créez un fichier nommé
producer.jset ajoutez-y le code suivant :const { MQClient, MessageProperties } = require('@aliyunmq/mq-http-sdk'); // HTTP endpoint of your ApsaraMQ for RocketMQ instance. const endpoint = "<http-endpoint>"; // Read credentials from environment variables. const accessKeyId = process.env['ALIBABA_CLOUD_ACCESS_KEY_ID']; const accessKeySecret = process.env['ALIBABA_CLOUD_ACCESS_KEY_SECRET']; const client = new MQClient(endpoint, accessKeyId, accessKeySecret); // Topic to publish to. Create this topic in the console first. const topic = "<topic>"; // Instance ID. Set to null or "" if the instance has no namespace. const instanceId = "<instance-id>"; const producer = client.getProducer(instanceId, topic); (async function () { try { for (let i = 0; i < 4; i++) { const msgProps = new MessageProperties(); // Custom property. msgProps.putProperty("a", i); // Message key for tracing or deduplication. msgProps.messageKey("MessageKey"); // Publish the message with a body, tag, and properties. const res = await producer.publishMessage("hello mq.", "TagA", msgProps); console.log("Publish message: MessageID:%s, BodyMD5:%s", res.body.MessageId, res.body.MessageBodyMD5); } } catch (e) { // Handle failures: retry or persist the message for later redelivery. console.log(e); } })(); -
Exécutez le producteur :
node producer.js -
Vérifiez la sortie. Une exécution réussie affiche quatre lignes similaires à celles-ci :
Publish message: MessageID:7F00000100246A4F0F2906B5BE4F0000, BodyMD5:4E69B245B8D0E47E5BE0B8520E8E608F Publish message: MessageID:7F00000100246A4F0F2906B5BE650001, BodyMD5:4E69B245B8D0E47E5BE0B8520E8E608F Publish message: MessageID:7F00000100246A4F0F2906B5BE780002, BodyMD5:4E69B245B8D0E47E5BE0B8520E8E608F Publish message: MessageID:7F00000100246A4F0F2906B5BE890003, BodyMD5:4E69B245B8D0E47E5BE0B8520E8E608F
Recevoir et acquitter des messages normaux
Le consommateur utilise le sondage long (long polling) : si aucun message n'est disponible, la requête reste bloquée sur le broker pendant un nombre de secondes configurable. Si un message arrive durant cette période, le broker répond immédiatement.
-
Créez un fichier nommé
consumer.jset ajoutez-y le code suivant :const { MQClient } = require('@aliyunmq/mq-http-sdk'); // HTTP endpoint of your ApsaraMQ for RocketMQ instance. const endpoint = "<http-endpoint>"; // Read credentials from environment variables. const accessKeyId = process.env['ALIBABA_CLOUD_ACCESS_KEY_ID']; const accessKeySecret = process.env['ALIBABA_CLOUD_ACCESS_KEY_SECRET']; const client = new MQClient(endpoint, accessKeyId, accessKeySecret); // Topic to consume from. const topic = "<topic>"; // Consumer group ID. Create this in the console first. const groupId = "<group-id>"; // Instance ID. Set to null or "" if the instance has no namespace. const instanceId = "<instance-id>"; const consumer = client.getConsumer(instanceId, topic, groupId); (async function () { while (true) { try { // consumeMessage(batchSize, pollingSeconds) // batchSize: max messages per request (1-16) // pollingSeconds: long-polling timeout in seconds (max 30) const res = await consumer.consumeMessage(3, 3); if (res.code === 200) { console.log("Consume messages, requestId: %s", res.requestId); const handles = res.body.map((message) => { console.log( "\tMessageId:%s, Tag:%s, PublishTime:%d, NextConsumeTime:%d, " + "FirstConsumeTime:%d, ConsumedTimes:%d, Body:%s, Props:%j, MessageKey:%s, Prop-A:%s", message.MessageId, message.MessageTag, message.PublishTime, message.NextConsumeTime, message.FirstConsumeTime, message.ConsumedTimes, message.MessageBody, message.Properties, message.MessageKey, message.Properties.a ); return message.ReceiptHandle; }); // Acknowledge consumed messages. If the broker does not receive // an ACK before NextConsumeTime, the message is delivered again. const ackRes = await consumer.ackMessage(handles); if (ackRes.code !== 204) { // If the handle of the message times out, the broker fails to receive an ACK. console.log("Ack failed:"); const failHandles = ackRes.body.map((error) => { console.log( "\tErrorHandle:%s, Code:%s, Reason:%s", error.ReceiptHandle, error.ErrorCode, error.ErrorMessage ); return error.ReceiptHandle; }); handles.forEach((handle) => { if (failHandles.indexOf(handle) < 0) { console.log("\tSucHandle:%s", handle); } }); } else { console.log("Ack succeeded, requestId: %s\n\t", ackRes.requestId, handles.join(',')); } } } catch (e) { if (e.Code && e.Code.indexOf("MessageNotExist") > -1) { // No messages available -- long polling continues on next iteration. console.log("No new messages. requestId: %s, Code: %s", e.RequestId, e.Code); } else { console.log(e); } } } })(); -
Exécutez le consommateur :
node consumer.js -
Vérifiez la sortie. Lorsque des messages sont disponibles, le consommateur affiche :
Consume messages, requestId: B1A341B6F2A6XXXX MessageId:7F00000100246A4F0F2906B5BE4F0000, Tag:TagA, PublishTime:1620000000000, NextConsumeTime:1620000030000, FirstConsumeTime:1620000000000, ConsumedTimes:1, Body:hello mq., Props:{"a":"0"}, MessageKey:MessageKey Ack succeeded, requestId: B1A341B6F2A6XXXXLorsqu'aucun message n'est disponible, il affiche :
No new messages. requestId: B1A341B6F2A6XXXX, Code: MessageNotExist
Méthodes API clés
| Méthode | Description |
|---|---|
new MQClient(endpoint, accessKeyId, accessKeySecret) |
Crée un client HTTP connecté à l'endpoint spécifié. |
client.getProducer(instanceId, topic) |
Renvoie un producteur lié à l'instance et au topic indiqués. |
client.getConsumer(instanceId, topic, groupId) |
Renvoie un consommateur lié à l'instance, au topic et au groupe de consommateurs indiqués. |
producer.publishMessage(body, tag, properties) |
Publie un message avec le corps, le tag et les propriétés facultatives spécifiés. La réponse contient MessageId et MessageBodyMD5. |
consumer.consumeMessage(batchSize, pollingSeconds) |
Récupère jusqu'à batchSize messages (maximum 16). Se bloque pendant un maximum de pollingSeconds secondes (maximum 30) si aucun message n'est disponible. |
consumer.ackMessage(handles) |
Acquitte les messages via leurs handles de réception. Renvoie le code d'état 204 en cas de succès. |
Acquittement des messages
Après avoir traité un message, acquittez-le en appelant ackMessage avec le handle de réception du message. Si le broker ne reçoit pas d'acquittement avant l'échéance NextConsumeTime, le message est redistribué.
À chaque consommation d'un message, le broker attribue un nouveau handle de réception comportant un horodatage unique. Utilisez toujours le dernier handle lors de l'acquittement.
Échecs d'acquittement courants :
| Scénario | Symptôme | Résolution |
|---|---|---|
| Expiration du handle | L'acquittement renvoie une erreur avec un code d'erreur indiquant l'expiration du handle | Traitez et acquittez les messages plus rapidement pour éviter l'expiration des handles. |
| Délai d'attente réseau dépassé | L'appel d'acquittement lève une exception | Réessayez l'acquittement avec le même handle avant son expiration. |
| Livraison en double | La valeur ConsumedTimes du message est supérieure à 1 |
Concevez un traitement idempotent des messages pour gérer les redistributions en toute sécurité. |