Cette rubrique propose des exemples de code C# pour envoyer et recevoir des messages planifiés et différés à l'aide du SDK client HTTP ApsaraMQ for RocketMQ.
Différence entre messages planifiés et messages différés
Les messages planifiés et les messages différés retardent tous deux la livraison, mais diffèrent par la manière dont vous spécifiez l'heure de livraison :
| Type | Comportement de livraison | Configuration de StartDeliverTime |
|---|---|---|
| Message différé | Livré après un délai fixe à compter de l'envoi | Horodatage actuel + durée du délai (en millisecondes) |
| Message planifié | Livré à un instant précis | Horodatage cible de livraison (en millisecondes) |
Via HTTP, les deux types utilisent la même propriété StartDeliverTime, qui correspond à un horodatage absolu en millisecondes. Le broker conserve le message jusqu'à cet horodatage, puis le livre aux consommateurs.
Pour plus d'informations, consultez la documentation Messages planifiés et messages différés.
Cas d'utilisation
Expiration des commandes : envoyez un message différé pour vérifier le statut du paiement après 30 minutes. Annulez la commande si elle n'est pas payée.
Nouvelle tentative avec backoff : après une défaillance transitoire, planifiez une nouvelle tentative avec un délai croissant.
Notifications programmées : envoyez un rappel à une date et heure futures spécifiques.
Prérequis
Avant de commencer, assurez-vous d'avoir :
Le SDK client HTTP C# installé. Pour les instructions de configuration, consultez Préparer l'environnement.
Une instance ApsaraMQ for RocketMQ, un topic et un groupe de consommateurs créés dans la console ApsaraMQ for RocketMQ.
Une paire AccessKey Alibaba Cloud. Pour les instructions, consultez Créer une paire AccessKey.
Envoyer des messages planifiés ou différés
La propriété StartDeliverTime détermine le moment où le broker livre chaque message. Définissez-la sur un horodatage Unix en millisecondes.
Pour un message différé, ajoutez la durée du délai à l'horodatage actuel :
AliyunSDKUtils.GetNowTimeStamp() + delayMs.Pour un message planifié, définissez le paramètre sur la différence de temps entre l'instant planifié et l'instant actuel.
Remplacez les espaces réservés suivants par vos valeurs réelles :
| Espace réservé | Description | Exemple |
|---|---|---|
<your-http-endpoint> |
Endpoint HTTP disponible sur la page Instance Details de la console ApsaraMQ for RocketMQ | http://xxxx.mqrest.cn-hangzhou.aliyuncs.com |
<your-topic> |
Nom du topic | scheduled-topic |
<your-instance-id> |
ID de l'instance (définissez-le sur null ou "" si l'instance ne possède pas de namespace) |
MQ_INST_xxxx |
using System;
using System.Collections.Generic;
using System.Threading;
using Aliyun.MQ.Model;
using Aliyun.MQ.Model.Exp;
using Aliyun.MQ.Util;
namespace Aliyun.MQ.Sample
{
public class ProducerSample
{
// HTTP endpoint. Get this from the Instance Details page in the ApsaraMQ for RocketMQ console.
private const string _endpoint = "<your-http-endpoint>";
// Get AccessKey credentials from environment variables.
// Do not hardcode credentials in your source code.
private const string _accessKeyId = Environment.GetEnvironmentVariable("ALIBABA_CLOUD_ACCESS_KEY_ID");
private const string _secretAccessKey = Environment.GetEnvironmentVariable("ALIBABA_CLOUD_ACCESS_KEY_SECRET");
// Topic to publish to. Create this in the ApsaraMQ for RocketMQ console.
private const string _topicName = "<your-topic>";
// Instance ID. If the instance has no namespace, set this to null or "".
// Check the Instance Details page to determine whether your instance has a namespace.
private const string _instanceId = "<your-instance-id>";
private static MQClient _client = new Aliyun.MQ.MQClient(_accessKeyId, _secretAccessKey, _endpoint);
static MQProducer producer = _client.GetProducer(_instanceId, _topicName);
static void Main(string[] args)
{
try
{
for (int i = 0; i < 4; i++)
{
TopicMessage sendMsg;
sendMsg = new TopicMessage("hello mq");
// Set a custom property on the message.
sendMsg.PutProperty("a", i.ToString());
// Key line: set the delivery time.
// Delayed message: add the delay (in ms) to the current timestamp.
// Scheduled message: set the parameter to the time difference between the scheduled point in time and the current point in time.
sendMsg.StartDeliverTime = AliyunSDKUtils.GetNowTimeStamp() + 10 * 1000; // 10-second delay
TopicMessage result = producer.PublishMessage(sendMsg);
Console.WriteLine("Message published: " + result);
}
}
catch (Exception ex)
{
Console.Write(ex);
}
}
}
}
Recevoir des messages planifiés ou différés
Consommez les messages planifiés et différés de la même manière que les messages normaux. Le broker les livre uniquement après l'horodatage StartDeliverTime ; aucune configuration spécifique du consommateur n'est donc requise.
L'exemple suivant utilise le long polling pour consommer les messages et acquitter chaque lot.
Remplacez les espaces réservés suivants par vos valeurs réelles :
| Espace réservé | Description | Exemple |
|---|---|---|
<your-http-endpoint> |
Endpoint HTTP disponible sur la page Instance Details | http://xxxx.mqrest.cn-hangzhou.aliyuncs.com |
<your-topic> |
Nom du topic | scheduled-topic |
<your-instance-id> |
ID de l'instance (définissez-le sur null ou "" si l'instance ne possède pas de namespace) |
MQ_INST_xxxx |
<your-group-id> |
ID du groupe de consommateurs | GID_scheduled_consumer |
using System;
using System.Collections.Generic;
using System.Threading;
using Aliyun.MQ.Model;
using Aliyun.MQ.Model.Exp;
using Aliyun.MQ;
namespace Aliyun.MQ.Sample
{
public class ConsumerSample
{
// HTTP endpoint. Get this from the Instance Details page in the ApsaraMQ for RocketMQ console.
private const string _endpoint = "<your-http-endpoint>";
// The AccessKey ID that is used for authentication.
private const string _accessKeyId = "${ACCESS_KEY}";
// The AccessKey secret that is used for authentication.
private const string _secretAccessKey = "${SECRET_KEY}";
// Topic to consume from. Create this in the ApsaraMQ for RocketMQ console.
private const string _topicName = "<your-topic>";
// Instance ID. If the instance has no namespace, set this to null or "".
private const string _instanceId = "<your-instance-id>";
// Consumer group ID. Create this in the ApsaraMQ for RocketMQ console.
private const string _groupId = "<your-group-id>";
private static MQClient _client = new Aliyun.MQ.MQClient(_accessKeyId, _secretAccessKey, _endpoint);
static MQConsumer consumer = _client.GetConsumer(_instanceId, _topicName, _groupId, null);
static void Main(string[] args)
{
// Consume messages in a loop. For production use, run multiple threads concurrently.
while (true)
{
try
{
// Long polling: if no messages are available, the request waits on the broker
// for up to the specified number of seconds before returning.
List<Message> messages = null;
try
{
messages = consumer.ConsumeMessage(
3, // Max messages per batch (up to 16)
3 // Long polling timeout in seconds (up to 30)
);
}
catch (Exception exp1)
{
if (exp1 is MessageNotExistException)
{
Console.WriteLine(Thread.CurrentThread.Name + " No new message, " + ((MessageNotExistException)exp1).RequestId);
continue;
}
Console.WriteLine(exp1);
Thread.Sleep(2000);
}
if (messages == null)
{
continue;
}
List<string> handlers = new List<string>();
Console.WriteLine(Thread.CurrentThread.Name + " Receive Messages:");
// Process each message.
foreach (Message message in messages)
{
Console.WriteLine(message);
Console.WriteLine("Property a is:" + message.GetProperty("a"));
handlers.Add(message.ReceiptHandle);
}
// Acknowledge messages. If the broker does not receive an ACK
// before Message.nextConsumeTime, it redelivers the message.
// Each delivery assigns a new receipt handle.
try
{
consumer.AckMessage(handlers);
Console.WriteLine("Ack message success:");
foreach (string handle in handlers)
{
Console.Write("\t" + handle);
}
Console.WriteLine();
}
catch (Exception exp2)
{
// An ACK can fail if the receipt handle has expired.
if (exp2 is AckMessageException)
{
AckMessageException ackExp = (AckMessageException)exp2;
Console.WriteLine("Ack message fail, RequestId:" + ackExp.RequestId);
foreach (AckMessageErrorItem errorItem in ackExp.ErrorItems)
{
Console.WriteLine("\tErrorHandle:" + errorItem.ReceiptHandle + ",ErrorCode:" + errorItem.ErrorCode + ",ErrorMsg:" + errorItem.ErrorMessage);
}
}
}
}
catch (Exception ex)
{
Console.WriteLine(ex);
Thread.Sleep(2000);
}
}
}
}
}
Voir aussi
Messages planifiés et messages différés – Concepts et limites
Créer des ressources – Configurer des instances, des topics et des groupes de consommateurs
Préparer l'environnement – Installer et configurer le SDK client HTTP C#