Tous les produits
Search
Centre de documentation

ApsaraMQ for RocketMQ:Scheduled and delayed messages

Dernière mise à jour :Aug 09, 2026

Les messages planifiés sont remis aux consommateurs à une heure précise plutôt qu'immédiatement. Utilisez-les pour mettre en œuvre des déclencheurs temporels dans vos applications distribuées, par exemple pour annuler une commande après expiration du délai, envoyer des notifications différées ou planifier des tâches périodiques.

Les messages planifiés et les messages différés fonctionnent de manière identique : le broker conserve les deux types et les remet à l'instant spécifié. Dans la suite de cette rubrique, le terme « messages planifiés » englobe ces deux catégories.

Prérequis

Avant d'envoyer des messages planifiés, vérifiez les points suivants :

  • Le paramètre MessageType de votre topic est défini sur Delay. L'envoi de messages planifiés n'est possible que vers des topics de type Delay.

  • L'horodatage de livraison est postérieur à l'heure actuelle. Si l'horodatage est dans le passé ou dépasse la plage autorisée, le broker remet le message immédiatement.

  • L'horodatage de livraison respecte la fenêtre de planification maximale associée à votre type d'instance :

    Type d'instance Fenêtre de planification maximale
    Édition Standard (abonnement, paiement à l'utilisation ou serverless) et Édition Professionnelle serverless 7 jours
    Édition Professionnelle et Édition Enterprise Platinum (abonnement ou paiement à l'utilisation) 40 jours

Cas d'utilisation

Orchestration distribuée des tâches

Planifiez l'envoi de messages à différentes granularités temporelles pour remplacer les minuteurs traditionnels basés sur une base de données. Par exemple, déclenchez quotidiennement une tâche de nettoyage de fichiers à heure fixe ou envoyez une notification de rappel toutes les 2 minutes. Les messages planifiés éliminent la complexité liée à la gestion des tâches cron distribuées tout en offrant une meilleure évolutivité.

Distributed timed scheduling

Gestion des délais d'expiration des tâches

Différez une action jusqu'à l'échéance d'un délai. Par exemple, dans le e-commerce, annulez une commande non payée 30 minutes après sa création au lieu d'interroger régulièrement une base de données pour vérifier le statut du paiement.

Task timeout processing

Par rapport à l'interrogation régulière d'une base de données pour identifier les enregistrements expirés, les messages planifiés offrent :

  • Une granularité temporelle flexible : déclenchez des tâches à n'importe quel intervalle, sans incréments de temps fixes ni besoin de logique de déduplication.

  • De meilleures performances : évitez les goulots d'étranglement liés aux analyses fréquentes de la base de données. L'infrastructure de messagerie gère la concurrence et monte en charge horizontalement.

Fonctionnement de la planification

Horodatage de livraison

L'heure de livraison est exprimée sous forme d'horodatage Unix en millisecondes. Définissez-la sur le message à l'aide de setDeliveryTimestamp(). Le broker conserve le message jusqu'à ce que cet horodatage soit atteint.

  • Pour les messages planifiés, calculez l'heure absolue de livraison : si l'heure actuelle est 2022-06-09 17:30:00 et que vous souhaitez une livraison à 19:20:00, définissez l'horodatage sur 1654773600000.

  • Pour les messages différés, ajoutez un décalage relatif à l'heure actuelle : si l'heure actuelle est 2022-06-09 17:30:00 et que vous souhaitez une livraison dans 1 heure, définissez l'horodatage sur 1654770600000.

Important

Il est impossible de modifier ou d'annuler l'horodatage de livraison après l'envoi du message.

Cycle de vie des messages

Un message planifié traverse six étapes, de sa production à sa suppression :

Scheduled message lifecycle

Étape Description
Initialisation Le producteur construit le message et l'envoie au broker.
Planification Le broker stocke le message dans un système de stockage basé sur le temps. Aucun index visible par les consommateurs n'est créé à ce stade.
Prêt pour la consommation À l'horodatage de livraison, le message est transféré vers le moteur de stockage standard et devient visible pour les consommateurs.
En cours de consommation Un consommateur reçoit et traite le message. Si aucun accusé de réception n'arrive avant l'expiration du délai, le broker retente la livraison. Pour plus d'informations, consultez Nouvelles tentatives de consommation.
Validé Le consommateur accuse réception du message. Le broker le marque comme consommé, mais ne le supprime pas immédiatement.
Suppression Le broker supprime le message de manière progressive lorsque la période de rétention expire ou lorsque l'espace de stockage vient à manquer. Jusqu'à la suppression, le message peut être re-consommé. Pour plus d'informations, consultez Stockage et nettoyage des messages.

Limites

Contrainte Détail
Type de topic Seuls les topics dont le paramètre MessageType est défini sur Delay acceptent les messages planifiés.
Précision temporelle Précision à la milliseconde près. La granularité par défaut est de 1 000 ms.
Persistance Les messages planifiés survivent aux redémarrages du broker. Toutefois, des exceptions du système de stockage ou des redémarrages peuvent entraîner de brefs retards de livraison.
Annulation Non prise en charge. Une fois envoyé, l'horodatage de livraison ne peut ni être modifié ni être révoqué.

Envoi et consommation de messages planifiés

Les exemples suivants utilisent le SDK Java Apache RocketMQ 5.x. Chaque exemple de producteur définit un horodatage de livraison avec setDeliveryTimestamp() ; c'est la seule différence par rapport à l'envoi d'un message normal. Pour la référence complète du SDK, consultez SDK Apache RocketMQ 5.x.

Assurez-vous que votre topic est créé avec le paramètre MessageType défini sur Delay avant d'exécuter ces exemples.

Exemple de code

Envoyer un message planifié

import org.apache.rocketmq.client.apis.*;
import org.apache.rocketmq.client.apis.message.Message;
import org.apache.rocketmq.client.apis.producer.Producer;
import org.apache.rocketmq.client.apis.producer.SendReceipt;

public class ProducerExample {
    public static void main(String[] args) throws ClientException {
        // Instance endpoint. Find this on the Endpoints tab of the Instance Details
        // page in the ApsaraMQ for RocketMQ console.
        // For internal access from ECS in the same VPC, use the VPC endpoint.
        // For public access, use the public endpoint and enable Internet access
        // on the instance.
        String endpoints = "rmq-cn-xxx.{regionId}.rmq.aliyuncs.com:8080";

        // Topic name. Create this topic in the console first, with MessageType
        // set to Delay.
        String topic = "Your Topic";

        ClientServiceProvider provider = ClientServiceProvider.loadService();
        ClientConfigurationBuilder builder = ClientConfiguration.newBuilder().setEndpoints(endpoints);

        // For public endpoint access, provide your instance username and password.
        // Find these on the Intelligent Authentication tab of the Access Control
        // page in the console.
        // For VPC access from ECS, credentials are obtained automatically.
        // For serverless instances with authentication-free in VPCs enabled,
        // credentials are not required for VPC access.
        //builder.setCredentialProvider(new StaticSessionCredentialsProvider("Instance UserName", "Instance Password"));

        ClientConfiguration configuration = builder.build();
        Producer producer = provider.newProducerBuilder()
                .setTopics(topic)
                .setClientConfiguration(configuration)
                .build();

        // Deliver the message 10 minutes from now
        long deliverTimeStamp = System.currentTimeMillis() + 10L * 60 * 1000;

        Message message = provider.newMessageBuilder()
                .setTopic("topic")
                .setKeys("messageKey")       // Message key for tracing
                .setTag("messageTag")        // Tag for consumer-side filtering
                .setDeliveryTimestamp(deliverTimeStamp)
                .setBody("messageBody".getBytes())
                .build();
        try {
            SendReceipt sendReceipt = producer.send(message);
            System.out.println(sendReceipt.getMessageId());
        } catch (ClientException e) {
            e.printStackTrace();
        }
    }
}

Consommer avec un consommateur push

Un consommateur push reçoit les messages via un callback de listener de messages. Le broker pousse les messages vers le consommateur dès qu'ils sont prêts.

import org.apache.rocketmq.client.apis.*;
import org.apache.rocketmq.client.apis.consumer.ConsumeResult;
import org.apache.rocketmq.client.apis.consumer.FilterExpression;
import org.apache.rocketmq.client.apis.consumer.FilterExpressionType;
import org.apache.rocketmq.client.apis.consumer.PushConsumer;
import org.apache.rocketmq.shaded.org.slf4j.Logger;
import org.apache.rocketmq.shaded.org.slf4j.LoggerFactory;

import java.io.IOException;
import java.util.Collections;

public class PushConsumerExample {
    private static final Logger LOGGER = LoggerFactory.getLogger(PushConsumerExample.class);

    private PushConsumerExample() {
    }

    public static void main(String[] args) throws ClientException, IOException, InterruptedException {
        // Instance endpoint (see producer example for details)
        String endpoints = "rmq-cn-xxx.{regionId}.rmq.aliyuncs.com:8080";

        // Create the topic and consumer group in the console before running
        // this example
        String topic = "Your Topic";
        String consumerGroup = "Your ConsumerGroup";

        final ClientServiceProvider provider = ClientServiceProvider.loadService();
        ClientConfigurationBuilder builder = ClientConfiguration.newBuilder().setEndpoints(endpoints);

        // Uncomment for public endpoint access (see producer example for details)
        //builder.setCredentialProvider(new StaticSessionCredentialsProvider("Instance UserName", "Instance Password"));
        ClientConfiguration clientConfiguration = builder.build();

        // Subscribe to all messages in the topic (tag = "*")
        String tag = "*";
        FilterExpression filterExpression = new FilterExpression(tag, FilterExpressionType.TAG);

        PushConsumer pushConsumer = provider.newPushConsumerBuilder()
                .setClientConfiguration(clientConfiguration)
                .setConsumerGroup(consumerGroup)
                .setSubscriptionExpressions(Collections.singletonMap(topic, filterExpression))
                .setMessageListener(messageView -> {
                    System.out.println("Consume Message: " + messageView);
                    return ConsumeResult.SUCCESS;
                })
                .build();

        Thread.sleep(Long.MAX_VALUE);
        // Shut down when no longer needed
        //pushConsumer.close();
    }
}

Consommer avec un consommateur simple

Un consommateur simple extrait (pull) les messages du broker et accuse explicitement réception de chacun d'eux après traitement.

import org.apache.rocketmq.client.apis.*;
import org.apache.rocketmq.client.apis.consumer.FilterExpression;
import org.apache.rocketmq.client.apis.consumer.FilterExpressionType;
import org.apache.rocketmq.client.apis.consumer.SimpleConsumer;
import org.apache.rocketmq.client.apis.message.MessageId;
import org.apache.rocketmq.client.apis.message.MessageView;
import org.apache.rocketmq.shaded.org.slf4j.Logger;
import org.apache.rocketmq.shaded.org.slf4j.LoggerFactory;

import java.io.IOException;
import java.time.Duration;
import java.util.Collections;
import java.util.List;

public class SimpleConsumerExample {
    private static final Logger LOGGER = LoggerFactory.getLogger(SimpleConsumerExample.class);

    private SimpleConsumerExample() {
    }

    public static void main(String[] args) throws ClientException, IOException {
        // Instance endpoint (see producer example for details)
        String endpoints = "rmq-cn-xxx.{regionId}.rmq.aliyuncs.com:8080";

        // Create the topic and consumer group in the console before running
        // this example
        String topic = "Your Topic";
        String consumerGroup = "Your ConsumerGroup";

        final ClientServiceProvider provider = ClientServiceProvider.loadService();
        ClientConfigurationBuilder builder = ClientConfiguration.newBuilder().setEndpoints(endpoints);

        // Uncomment for public endpoint access (see producer example for details)
        //builder.setCredentialProvider(new StaticSessionCredentialsProvider("Instance UserName", "Instance Password"));
        ClientConfiguration clientConfiguration = builder.build();

        // Long-polling timeout
        Duration awaitDuration = Duration.ofSeconds(10);

        // Subscribe to all messages in the topic (tag = "*")
        String tag = "*";
        FilterExpression filterExpression = new FilterExpression(tag, FilterExpressionType.TAG);

        SimpleConsumer consumer = provider.newSimpleConsumerBuilder()
                .setClientConfiguration(clientConfiguration)
                .setConsumerGroup(consumerGroup)
                .setAwaitDuration(awaitDuration)
                .setSubscriptionExpressions(Collections.singletonMap(topic, filterExpression))
                .build();

        int maxMessageNum = 16;
        Duration invisibleDuration = Duration.ofSeconds(10);

        // Pull and process messages in a loop.
        // For real-time consumption, use multiple threads to pull concurrently.
        while (true) {
            final List<MessageView> messages = consumer.receive(maxMessageNum, invisibleDuration);
            messages.forEach(messageView -> {
                System.out.println("Received message: " + messageView);
            });
            for (MessageView message : messages) {
                final MessageId messageId = message.getMessageId();
                try {
                    consumer.ack(message);
                    System.out.println("Message is acknowledged successfully, messageId= " + messageId);
                } catch (Throwable t) {
                    t.printStackTrace();
                }
            }
        }
        // Shut down when no longer needed
        // consumer.close();
    }
}

Bonnes pratiques

Échelonnez les heures de livraison. Évitez de planifier un grand lot de messages pour exactement le même horodatage. Lorsque de nombreux messages partagent la même heure de livraison, le broker doit tous les traiter simultanément, ce qui augmente la charge et peut retarder la remise. Répartissez les horodatages sur une plage lorsque cela est possible.

Intégrez une logique d'annulation côté consommateur. Puisqu'il est impossible de révoquer ou de replanifier un message après son envoi, gérez l'annulation au niveau du consommateur. Par exemple, lorsqu'un consommateur reçoit un message « annuler la commande non payée », vérifiez d'abord le statut actuel de la commande ; si l'utilisateur a déjà payé, ignorez l'annulation.

FAQ

Puis-je annuler ou replanifier un message après l'avoir envoyé ?

Non. L'horodatage de livraison est définitif une fois le message envoyé. Pour gérer les scénarios où une action planifiée devient inutile, intégrez cette logique dans votre consommateur. Par exemple, vérifiez le statut actuel de la commande avant d'annuler une commande non payée.

Que se passe-t-il si je définis une heure de livraison dans le passé ?

Le broker le remet immédiatement, comme un message normal.

Pourquoi ne puis-je pas trouver mon message planifié dans la console ?

Les messages planifiés restent invisibles jusqu'à ce que l'horodatage de livraison soit atteint. Vérifiez après cet instant : le message apparaît alors et devient disponible pour les consommateurs.