Tous les produits
Search
Centre de documentation

ApsaraMQ for RocketMQ:Send and receive scheduled and delayed messages over HTTP (C#)

Dernière mise à jour :Aug 09, 2026

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 :

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