Tous les produits
Search
Centre de documentation

ApsaraMQ for RocketMQ:Send and consume scheduled or delayed messages

Dernière mise à jour :Aug 09, 2026

Cette rubrique propose des exemples de code Java pour envoyer et consommer des messages planifiés ou différés via le SDK client HTTP d'ApsaraMQ for RocketMQ.

Messages planifiés et messages différés

Au niveau de l'API, les messages planifiés et les messages différés constituent une seule et même fonctionnalité. Tous deux utilisent la méthode setStartDeliverTime pour différer la livraison. La seule différence réside dans le calcul de l'horodatage :

Type Méthode de définition de l'heure de livraison Exemple
Message différé Heure actuelle + décalage temporel System.currentTimeMillis() + 10 * 1000 (délai de 10 secondes)
Message planifié Horodatage Unix absolu en millisecondes 1654770600000 (2022-06-09 18:30:00)

Le broker conserve chaque message jusqu'à l'heure de livraison spécifiée, puis le libère pour consommation. Ce paramètre s'applique par message, ce qui permet à chaque message d'avoir sa propre heure de livraison.

Pour plus d'informations, consultez la documentation Messages planifiés et messages différés.

Prérequis

Avant de commencer, assurez-vous d'avoir :

Important

Chaque topic ne prend en charge qu'un seul type de message. N'utilisez pas un topic configuré pour des messages normaux pour envoyer des messages planifiés ou différés.

Envoyer des messages planifiés ou différés

Tous les messages de cet exemple utilisent la méthode setStartDeliverTime pour différer la livraison de 10 secondes. Pour envoyer un message planifié, remplacez le décalage relatif par un horodatage Unix absolu en millisecondes (par exemple, 1654770600000 pour le 09/06/2022 à 18:30:00).

import com.aliyun.mq.http.MQClient;
import com.aliyun.mq.http.MQProducer;
import com.aliyun.mq.http.model.TopicMessage;

import java.util.Date;

public class Producer {

    public static void main(String[] args) {
        MQClient mqClient = new MQClient(
                // HTTP endpoint. Find this on the Instance Details page, in the
                // Endpoints tab under Basic Information.
                "<your-http-endpoint>",
                // Load credentials from environment variables to avoid hardcoding secrets.
                System.getenv("ALIBABA_CLOUD_ACCESS_KEY_ID"),
                System.getenv("ALIBABA_CLOUD_ACCESS_KEY_SECRET")
        );

        final String topic = "<your-topic>";
        // If the instance has a namespace, specify the instance ID.
        // If it does not have a namespace, set this to null or "".
        final String instanceId = "<your-instance-id>";

        MQProducer producer;
        if (instanceId != null && instanceId != "") {
            producer = mqClient.getProducer(instanceId, topic);
        } else {
            producer = mqClient.getProducer(topic);
        }

        try {
            for (int i = 0; i < 4; i++) {
                TopicMessage pubMsg = new TopicMessage(
                        "hello mq!".getBytes(),   // Message body
                        "A"                        // Message tag
                );
                pubMsg.getProperties().put("a", String.valueOf(i));  // Custom property
                pubMsg.setMessageKey("MessageKey");

                // Defer delivery by 10 seconds.
                // For a scheduled message, use an absolute timestamp instead,
                // e.g., pubMsg.setStartDeliverTime(1654770600000L);
                pubMsg.setStartDeliverTime(System.currentTimeMillis() + 10 * 1000);

                TopicMessage pubResultMsg = producer.publishMessage(pubMsg);

                System.out.println(new Date() + " Send mq message success."
                        + " Topic is:" + topic
                        + ", msgId is: " + pubResultMsg.getMessageId()
                        + ", bodyMD5 is: " + pubResultMsg.getMessageBodyMD5());
            }
        } catch (Throwable e) {
            System.out.println(new Date() + " Send mq message failed. Topic is:" + topic);
            e.printStackTrace();
        }

        mqClient.close();
    }
}

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

Espace réservé Description Emplacement
<your-http-endpoint> Endpoint HTTP Page Instance Details > onglet Endpoints > HTTP Endpoint
<your-topic> Nom du topic Console ApsaraMQ for RocketMQ > Topics
<your-instance-id> ID de l'instance, ou null si l'instance n'a pas de namespace Page Instance Details > section Basic Information

Consommer 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 conserve jusqu'à l'heure de livraison, puis les libère pour consommation.

L'exemple suivant utilise le long polling pour consommer les messages et renvoyer des accusés de réception (ACK) au broker.

import com.aliyun.mq.http.MQClient;
import com.aliyun.mq.http.MQConsumer;
import com.aliyun.mq.http.common.AckMessageException;
import com.aliyun.mq.http.model.Message;

import java.util.ArrayList;
import java.util.List;

public class Consumer {

    public static void main(String[] args) {
        MQClient mqClient = new MQClient(
                // HTTP endpoint. Find this on the Instance Details page, in the
                // Endpoints tab under Basic Information.
                "<your-http-endpoint>",
                // The AccessKey ID for authentication.
                "${ACCESS_KEY}",
                // The AccessKey secret for authentication.
                "${SECRET_KEY}"
        );

        final String topic = "<your-topic>";
        final String groupId = "<your-group-id>";
        // If the instance has a namespace, specify the instance ID.
        // If it does not have a namespace, set this to null or "".
        final String instanceId = "<your-instance-id>";

        final MQConsumer consumer;
        if (instanceId != null && instanceId != "") {
            consumer = mqClient.getConsumer(instanceId, topic, groupId, null);
        } else {
            consumer = mqClient.getConsumer(topic, groupId);
        }

        // Consume messages in a loop. For production use, run multiple threads
        // to consume messages concurrently.
        do {
            List<Message> messages = null;

            try {
                // Long polling: wait up to 3 seconds for new messages.
                // - First parameter: max messages per batch (max: 16).
                // - Second parameter: long polling timeout in seconds (max: 30).
                messages = consumer.consumeMessage(3, 3);
            } catch (Throwable e) {
                e.printStackTrace();
                try {
                    Thread.sleep(2000);
                } catch (InterruptedException e1) {
                    e1.printStackTrace();
                }
            }

            if (messages == null || messages.isEmpty()) {
                System.out.println(Thread.currentThread().getName()
                        + ": no new message, continue!");
                continue;
            }

            // Process messages
            for (Message message : messages) {
                System.out.println("Receive message: " + message);
            }

            // ACK consumed messages. If the broker does not receive an ACK
            // before the retry interval elapses, the message is redelivered.
            {
                List<String> handles = new ArrayList<String>();
                for (Message message : messages) {
                    handles.add(message.getReceiptHandle());
                }

                try {
                    consumer.ackMessage(handles);
                } catch (Throwable e) {
                    if (e instanceof AckMessageException) {
                        AckMessageException errors = (AckMessageException) e;
                        System.out.println("Ack message fail, requestId is:"
                                + errors.getRequestId() + ", fail handles:");
                        if (errors.getErrorMessages() != null) {
                            for (String errorHandle : errors.getErrorMessages().keySet()) {
                                System.out.println("Handle:" + errorHandle
                                        + ", ErrorCode:" + errors.getErrorMessages().get(errorHandle).getErrorCode()
                                        + ", ErrorMsg:" + errors.getErrorMessages().get(errorHandle).getErrorMessage());
                            }
                        }
                        continue;
                    }
                    e.printStackTrace();
                }
            }
        } while (true);
    }
}

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

Espace réservé Description Emplacement
<your-http-endpoint> Endpoint HTTP Page Instance Details > onglet Endpoints > HTTP Endpoint
${ACCESS_KEY} AccessKey ID pour l'authentification Consultez Create an AccessKey pair
${SECRET_KEY} AccessKey secret pour l'authentification Consultez Create an AccessKey pair
<your-topic> Nom du topic Console ApsaraMQ for RocketMQ > Topics
<your-group-id> ID du groupe de consommateurs Console ApsaraMQ for RocketMQ > Groups
<your-instance-id> ID de l'instance, ou null si l'instance n'a pas de namespace Page Instance Details > section Basic Information

Étapes suivantes