As mensagens atrasadas permitem controlar quando uma mensagem fica disponível para consumo. O broker do ApsaraMQ for RocketMQ retém a mensagem atrasada até o horário de entrega especificado e, em seguida, a libera para os assinantes. Esse mecanismo funciona de forma semelhante às filas com atraso.
Casos de uso comuns:
Tratamento de timeout de pedidos: cancele um pedido não pago após 30 minutos ao enviar uma mensagem atrasada que acione a lógica de cancelamento.
Nova tentativa com backoff: agende uma nova tentativa após um atraso fixo quando uma dependência upstream estiver temporariamente indisponível.
Notificações diferidas: envie um lembrete ou notificação de acompanhamento após um período de espera configurável.
Para conceitos e restrições, consulte Mensagens agendadas e mensagens atrasadas.
Pré-requisitos
Antes de começar, verifique se você tem:
O SDK de cliente TCP para Java instalado. Para mais detalhes, consulte Preparar o ambiente.
Uma instância do ApsaraMQ for RocketMQ, um tópico e um grupo de consumidores criados no console. Para mais detalhes, consulte Criar recursos.
Um par de AccessKey da sua conta Alibaba Cloud. Para mais detalhes, consulte Criar um par de AccessKey.
(Opcional) Logging configurado. Para mais detalhes, consulte Configurações de logging.
Como funcionam as mensagens atrasadas
O produtor cria uma mensagem e define um horário de entrega usando
setStartDeliverTime().O broker recebe a mensagem e a retém até o horário de entrega especificado.
No horário de entrega, o broker libera a mensagem no tópico e a torna disponível para os consumidores.
|
Parâmetro |
Detalhe |
|
Formato de hora |
Timestamp absoluto em milissegundos (Unix epoch) |
|
Valor mínimo |
Deve ser posterior à hora atual |
|
Atraso máximo |
40 dias a partir da hora atual |
Enviar mensagens atrasadas
A única diferença em relação ao envio de uma mensagem normal é uma chamada de método: msg.setStartDeliverTime(delayTime). Isso define o timestamp absoluto (em milissegundos) no qual o broker entrega a mensagem.
Para mais exemplos, consulte a Biblioteca de código do ApsaraMQ for RocketMQ.
import com.aliyun.openservices.ons.api.Message;
import com.aliyun.openservices.ons.api.ONSFactory;
import com.aliyun.openservices.ons.api.Producer;
import com.aliyun.openservices.ons.api.PropertyKeyConst;
import com.aliyun.openservices.ons.api.SendResult;
import java.util.Date;
import java.util.Properties;
public class ProducerDelayTest {
public static void main(String[] args) {
Properties properties = new Properties();
// Retrieve the AccessKey pair from environment variables.
// Make sure ALIBABA_CLOUD_ACCESS_KEY_ID and ALIBABA_CLOUD_ACCESS_KEY_SECRET are set.
properties.put(PropertyKeyConst.AccessKey, System.getenv("ALIBABA_CLOUD_ACCESS_KEY_ID"));
properties.put(PropertyKeyConst.SecretKey, System.getenv("ALIBABA_CLOUD_ACCESS_KEY_SECRET"));
// Specify the TCP endpoint.
// Find this value in the TCP Endpoint section of the Instance Details page
// in the ApsaraMQ for RocketMQ console.
properties.put(PropertyKeyConst.NAMESRV_ADDR, "<your-tcp-endpoint>");
Producer producer = ONSFactory.createProducer(properties);
// Call start() once before sending any messages.
producer.start();
Message msg = new Message(
"<your-topic>", // Topic created in the console
"DelayMessageTag", // Tag for consumer-side filtering, similar to a Gmail tag
"Hello MQ".getBytes() // Message body in binary format.
// ApsaraMQ for RocketMQ does not process message bodies.
// Producer and consumer must agree on serialization.
);
// Optional: set a business key for message tracing.
// The key should be globally unique when possible. Use it to query
// or resend messages in the ApsaraMQ for RocketMQ console.
msg.setKey("ORDERID_100");
try {
// Set the delivery time to 3 seconds from now.
// The value is an absolute timestamp in milliseconds.
// The maximum delay that you can specify is 40 days.
long delayTime = System.currentTimeMillis() + 3000;
msg.setStartDeliverTime(delayTime);
// Send the message in synchronous transmission mode.
// If no exception is thrown, the message is sent.
SendResult sendResult = producer.send(msg);
if (sendResult != null) {
System.out.println(new Date() + " Send mq message success. Topic is:"
+ msg.getTopic() + " msgId is: " + sendResult.getMessageId());
}
} catch (Exception e) {
// Handle send failure: log, retry, or persist the message.
System.out.println(new Date() + " Send mq message failed. Topic is:" + msg.getTopic());
e.printStackTrace();
}
// Shut down the producer when the application exits.
// For high-throughput scenarios, keep the producer running and reuse it.
producer.shutdown();
}
}
Substitua os seguintes espaços reservados pelos valores reais:
|
Espaço reservado |
Descrição |
Exemplo |
|
|
Endpoint TCP da página Detalhes da Instância |
|
|
|
Nome do tópico criado no console |
|
Se você está começando agora com o ApsaraMQ for RocketMQ, consulte o Projeto de demonstração para configurar um projeto funcional antes de enviar e receber mensagens.
Inscrever-se em mensagens atrasadas
A inscrição em mensagens atrasadas funciona da mesma forma que a inscrição em mensagens normais. Nenhuma configuração especial do consumidor é necessária, pois a lógica de atraso reside inteiramente no lado do produtor.
Para o código de exemplo do consumidor, consulte Inscrever-se em mensagens.
Próximos passos
Mensagens agendadas e mensagens atrasadas: compreenda os conceitos e as restrições da entrega de mensagens atrasadas e agendadas.