Tous les produits
Search
Centre de documentation

ApsaraMQ for RocketMQ:Send and receive scheduled and delayed messages

Dernière mise à jour :Aug 09, 2026

ApsaraMQ for RocketMQ prend en charge deux types de livraison de messages basée sur le temps : les messages planifiés, livrés à un instant précis, et les messages différés, livrés après un délai fixe. Ces deux types utilisent la propriété __STARTDELIVERTIME pour contrôler le moment où le broker livre le message aux consommateurs.

L'exemple de code Java ci-dessous montre comment envoyer des messages planifiés et différés, ainsi que s'y abonner, avec le SDK client TCP (Community Edition).

Fonctionnement

Lorsqu'un producteur envoie un message avec la propriété __STARTDELIVERTIME, le broker ApsaraMQ for RocketMQ conserve le message jusqu'à l'heure de livraison spécifiée. À ce moment-là, le broker livre le message aux consommateurs abonnés comme n'importe quel message normal.

  • Message différé : Définissez __STARTDELIVERTIME sur System.currentTimeMillis() + delayMillis. Le broker livre le message une fois le délai écoulé.

  • Message planifié : Définissez __STARTDELIVERTIME sur un horodatage Unix futur en millisecondes. Le broker livre le message exactement à cet instant.

Si l'heure spécifiée est dans le passé, le broker livre immédiatement le message.

Planifié ou différé : quand utiliser chacun

Type Cas d'utilisation Exemple
Différé Déclencher une action après une attente fixe Annuler une commande impayée 30 minutes après sa création
Planifié Déclencher une action à une heure précise Envoyer une notification à 09:00 chaque lundi

Différences par rapport à Apache RocketMQ

Apache RocketMQ prend en charge les messages différés, mais pas les messages planifiés ; aucune interface dédiée n'existe pour ces derniers.

ApsaraMQ for RocketMQ offre des fonctionnalités supplémentaires :

  • Prise en charge des messages différés et planifiés

  • Précision à la seconde près pour l'heure de livraison

  • Concurrence plus élevée pour le traitement des messages basés sur le temps

Important

Les méthodes de configuration et les résultats diffèrent entre Apache RocketMQ et ApsaraMQ for RocketMQ. Utilisez l'exemple de code de cette page pour ApsaraMQ for RocketMQ sur le cloud.

Prérequis

Avant de commencer, assurez-vous d'avoir :

Envoi de messages planifiés et différés

Les messages planifiés et différés utilisent tous deux la même propriété utilisateur __STARTDELIVERTIME. La seule différence réside dans le calcul de la valeur de l'horodatage.

Remplacez les espaces réservés suivants par vos valeurs réelles :

Espace réservé Description Exemple
<your-group-id> ID de groupe créé dans la console ApsaraMQ for RocketMQ GID_example
<your-access-point> Endpoint de l'instance depuis la console http://MQ_INST_XXXX.aliyuncs.com:80
<your-topic> Topic créé dans la console ApsaraMQ for RocketMQ Topic_example
<your-message-tag> Tag de message pour le filtrage TagA
import java.text.SimpleDateFormat;
import java.util.Date;
import org.apache.rocketmq.acl.common.AclClientRPCHook;
import org.apache.rocketmq.acl.common.SessionCredentials;
import org.apache.rocketmq.client.AccessChannel;
import org.apache.rocketmq.client.exception.MQClientException;
import org.apache.rocketmq.client.producer.DefaultMQProducer;
import org.apache.rocketmq.client.producer.SendResult;
import org.apache.rocketmq.common.message.Message;
import org.apache.rocketmq.remoting.RPCHook;
import org.apache.rocketmq.remoting.common.RemotingHelper;

public class RocketMQProducer {
    /**
     * Reads credentials from environment variables:
     *   ALIBABA_CLOUD_ACCESS_KEY_ID
     *   ALIBABA_CLOUD_ACCESS_KEY_SECRET
     */
    private static RPCHook getAclRPCHook() {
        return new AclClientRPCHook(new SessionCredentials(
            System.getenv("ALIBABA_CLOUD_ACCESS_KEY_ID"),
            System.getenv("ALIBABA_CLOUD_ACCESS_KEY_SECRET")));
    }

    public static void main(String[] args) throws MQClientException {
        // Create a producer with message trace enabled.
        // To disable message trace, use:
        //   new DefaultMQProducer("<your-group-id>", getAclRPCHook());
        DefaultMQProducer producer = new DefaultMQProducer(
            "<your-group-id>", getAclRPCHook(), true, null);

        // Required for message trace on the cloud.
        producer.setAccessChannel(AccessChannel.CLOUD);

        // Set the instance endpoint from the ApsaraMQ for RocketMQ console.
        producer.setNamesrvAddr("<your-access-point>");
        producer.start();

        for (int i = 0; i < 128; i++) {
            try {
                Message msg = new Message(
                    "<your-topic>",
                    "<your-message-tag>",
                    "Hello world".getBytes(RemotingHelper.DEFAULT_CHARSET));

                // --- Option A: Delayed message ---
                // Deliver 3 seconds from now.
                long delayTime = System.currentTimeMillis() + 3000;
                msg.putUserProperty("__STARTDELIVERTIME", String.valueOf(delayTime));

                // --- Option B: Scheduled message ---
                // Deliver at a specific date and time.
                // Uncomment the following lines and comment out Option A to use.
                //
                // long timeStamp = new SimpleDateFormat("yyyy-MM-dd HH:mm:ss")
                //     .parse("2025-08-10 18:45:00").getTime();
                // msg.putUserProperty("__STARTDELIVERTIME", String.valueOf(timeStamp));

                SendResult sendResult = producer.send(msg);
                System.out.printf("%s%n", sendResult);
            } catch (Exception e) {
                System.out.println(new Date() + " Send mq message failed.");
                e.printStackTrace();
            }
        }

        // Shut down the producer before exiting (optional).
        producer.shutdown();
    }
}

Résumé :

  • La propriété __STARTDELIVERTIME accepte un horodatage Unix en millisecondes.

  • Pour les messages différés, ajoutez le délai souhaité (en millisecondes) à l'heure actuelle.

  • Pour les messages planifiés, convertissez la chaîne de date-heure cible au format yyyy-MM-dd HH:mm:ss en horodatage Unix.

  • Si l'heure spécifiée est antérieure à l'heure actuelle, le message est livré immédiatement.

Abonnement aux messages planifiés et différés

L'abonnement aux messages planifiés et différés est identique à celui des messages normaux. Le broker conserve le message jusqu'à l'heure de livraison spécifiée, donc aucune configuration spéciale côté consommateur n'est nécessaire.

import java.util.List;
import org.apache.rocketmq.acl.common.AclClientRPCHook;
import org.apache.rocketmq.acl.common.SessionCredentials;
import org.apache.rocketmq.client.AccessChannel;
import org.apache.rocketmq.client.consumer.DefaultMQPushConsumer;
import org.apache.rocketmq.client.consumer.listener.ConsumeConcurrentlyContext;
import org.apache.rocketmq.client.consumer.listener.ConsumeConcurrentlyStatus;
import org.apache.rocketmq.client.consumer.listener.MessageListenerConcurrently;
import org.apache.rocketmq.client.consumer.rebalance.AllocateMessageQueueAveragely;
import org.apache.rocketmq.client.exception.MQClientException;
import org.apache.rocketmq.common.message.MessageExt;
import org.apache.rocketmq.remoting.RPCHook;

public class RocketMQPushConsumer {
    /**
     * Reads credentials from environment variables:
     *   ALIBABA_CLOUD_ACCESS_KEY_ID
     *   ALIBABA_CLOUD_ACCESS_KEY_SECRET
     */
    private static RPCHook getAclRPCHook() {
        return new AclClientRPCHook(new SessionCredentials(
            System.getenv("ALIBABA_CLOUD_ACCESS_KEY_ID"),
            System.getenv("ALIBABA_CLOUD_ACCESS_KEY_SECRET")));
    }

    public static void main(String[] args) throws MQClientException {
        // Create a consumer with message trace enabled.
        // To disable message trace, use:
        //   new DefaultMQPushConsumer("<your-group-id>", getAclRPCHook(),
        //       new AllocateMessageQueueAveragely());
        DefaultMQPushConsumer consumer = new DefaultMQPushConsumer(
            "<your-group-id>", getAclRPCHook(),
            new AllocateMessageQueueAveragely(), true, null);

        // Set the instance endpoint from the ApsaraMQ for RocketMQ console.
        // The value is in the format of http://xxxx.mq-internet.aliyuncs.com:80.
        consumer.setNamesrvAddr("<your-access-point>");

        // Required for message trace on the cloud.
        consumer.setAccessChannel(AccessChannel.CLOUD);

        // Subscribe to all tags (*) on the topic.
        consumer.subscribe("<your-topic>", "*");

        consumer.registerMessageListener(new MessageListenerConcurrently() {
            @Override
            public ConsumeConcurrentlyStatus consumeMessage(
                    List<MessageExt> msgs,
                    ConsumeConcurrentlyContext context) {
                System.out.printf("Receive New Messages: %s %n", msgs);
                return ConsumeConcurrentlyStatus.CONSUME_SUCCESS;
            }
        });

        consumer.start();
    }
}

Voir aussi