Todos os produtos
Search
Central de documentação

ApsaraMQ for RocketMQ:Scheduled and delayed messages

Última atualização: Jun 27, 2026

Mensagens agendadas são entregues aos consumidores em um horário específico, em vez de imediatamente. Use-as para criar gatilhos baseados em tempo em aplicações distribuídas, como cancelamento de pedidos por timeout, notificações adiadas ou agendamento periódico de tarefas.

Mensagens agendadas e mensagens atrasadas funcionam da mesma forma: o broker retém ambas e as entrega no timestamp especificado. O termo "mensagens agendadas" a seguir abrange os dois tipos.

Antes de começar

Antes de enviar mensagens agendadas, verifique se:

  • O tópico tem o MessageType definido como Delay. Mensagens agendadas só podem ser enviadas para tópicos do tipo Delay.

  • O timestamp de entrega é posterior ao horário atual. Se o timestamp estiver no passado ou exceder o intervalo permitido, o broker entrega a mensagem imediatamente.

  • O timestamp de entrega está dentro da janela máxima de agendamento para o tipo da sua instância:

    Tipo de instância

    Janela máxima de agendamento

    Standard Edition (assinatura, pagamento conforme o uso ou serverless) e Professional Edition serverless

    7 dias

    Professional Edition e Enterprise Platinum Edition (assinatura ou pagamento conforme o uso)

    40 dias

Casos de uso

Agendamento de tarefas distribuídas

Agende mensagens com diferentes granularidades de tempo para substituir temporizadores tradicionais baseados em banco de dados. Por exemplo, acione uma tarefa diária de limpeza de arquivos em um horário fixo ou envie uma notificação de lembrete a cada 2 minutos. As mensagens agendadas eliminam a complexidade de gerenciar cron jobs distribuídos e oferecem maior escalabilidade.

Distributed timed scheduling

Tratamento de timeout de tarefas

Adie uma ação até que um prazo expire. No comércio eletrônico, por exemplo, cancele um pedido não pago 30 minutos após a criação, em vez de consultar o banco de dados repetidamente para verificar o status do pagamento.

Task timeout processing

Em comparação com a consulta frequente ao banco de dados para buscar registros expirados, as mensagens agendadas oferecem:

  • Granularidade de tempo flexível: acione tarefas em qualquer intervalo, sem incrementos de tempo fixos e sem necessidade de lógica de deduplicação.

  • Maior desempenho: evite gargalos causados por varreduras frequentes no banco de dados. A infraestrutura de mensageria gerencia a concorrência e escala horizontalmente.

Como funciona o agendamento

Timestamp de entrega

O horário de entrega é expresso como um timestamp Unix em milissegundos. Defina-o na mensagem com setDeliveryTimestamp(). O broker retém a mensagem até que esse timestamp seja atingido.

  • Para mensagens agendadas, calcule o horário absoluto de entrega: se o horário atual for 2022-06-09 17:30:00 e você desejar a entrega às 19:20:00, defina o timestamp como 1654773600000.

  • Para mensagens atrasadas, adicione um deslocamento relativo ao horário atual: se o horário atual for 2022-06-09 17:30:00 e você desejar a entrega após 1 hora, defina o timestamp como 1654770600000.

Importante

Não é possível alterar ou cancelar o timestamp de entrega após o envio da mensagem.

Ciclo de vida da mensagem

Uma mensagem agendada passa por seis estágios, desde a produção até a exclusão:

Scheduled message lifecycle

Estágio

O que acontece

Inicialização

O produtor constrói a mensagem e a envia ao broker.

Agendada

O broker armazena a mensagem em um sistema de armazenamento baseado em tempo. Nenhum índice visível para o consumidor é criado neste momento.

Pronta para consumo

No timestamp de entrega, a mensagem é movida para o mecanismo de armazenamento regular e torna-se visível para os consumidores.

Em consumo

Um consumidor recebe e processa a mensagem. Se nenhum reconhecimento chegar dentro do timeout, o broker tenta entregar novamente. Para mais informações, consulte Tentativa de consumo.

Confirmada

O consumidor reconhece a mensagem. O broker a marca como consumida, mas não a exclui imediatamente.

Excluída

O broker exclui a mensagem de forma rotativa quando o período de retenção expira ou o espaço de armazenamento fica escasso. Até a exclusão, a mensagem pode ser reconsumida. Para mais informações, consulte Armazenamento e limpeza de mensagens.

Limites

Restrição

Detalhe

Tipo de tópico

Apenas tópicos com MessageType definido como Delay aceitam mensagens agendadas.

Precisão de tempo

Precisão de milissegundos. A granularidade padrão é de 1.000 ms.

Persistência

Mensagens agendadas sobrevivem a reinicializações do broker. No entanto, exceções ou reinicializações do sistema de armazenamento podem causar breves atrasos na entrega.

Cancelamento

Não suportado. Após o envio, não é possível alterar ou revogar o timestamp de entrega.

Enviar e consumir mensagens agendadas

Os exemplos abaixo usam o SDK Java do Apache RocketMQ 5.x. Cada exemplo de produtor define um timestamp de entrega com setDeliveryTimestamp(); essa é a única diferença em relação ao envio de uma mensagem normal. Para a referência completa do SDK, consulte SDKs do Apache RocketMQ 5.x.

Certifique-se de que seu tópico foi criado com MessageType definido como Delay antes de executar estes exemplos.

Código de exemplo

Enviar uma mensagem agendada

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();
        }
    }
}

Consumir com um push consumer

Um push consumer recebe mensagens por meio de um callback de listener de mensagens. O broker envia as mensagens para o consumidor assim que elas ficam prontas.

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();
    }
}

Consumir com um simple consumer

Um simple consumer busca mensagens do broker e reconhece cada uma explicitamente após o processamento.

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();
    }
}

Melhores práticas

Distribua os horários de entrega. Evite agendar um grande lote de mensagens para exatamente o mesmo timestamp. Quando muitas mensagens compartilham o mesmo horário de entrega, o broker precisa processar todas simultaneamente, o que aumenta a carga e pode atrasar a entrega. Sempre que possível, espalhe os timestamps por um intervalo de tempo.

Implemente lógica de cancelamento nos consumidores. Como não é possível revogar ou reagendar uma mensagem após o envio, trate o cancelamento no lado do consumidor. Por exemplo, ao receber uma mensagem de "cancelar pedido não pago", verifique primeiro o status atual do pedido; se o usuário já tiver pago, ignore o cancelamento.

Perguntas frequentes

Posso cancelar ou reagendar uma mensagem após o envio?

Não. O timestamp de entrega é definitivo após o envio da mensagem. Para lidar com cenários em que uma ação agendada se torna desnecessária, implemente essa lógica no consumidor. Por exemplo, verifique o status atual do pedido antes de cancelar um pedido não pago.

O que acontece se eu definir um horário de entrega no passado?

O broker a entrega imediatamente, assim como uma mensagem normal.

Por que não consigo encontrar minha mensagem agendada no console?

Mensagens agendadas permanecem invisíveis até que o timestamp de entrega seja atingido. Verifique novamente após esse horário; a mensagem aparecerá e estará disponível para os consumidores.