O ApsaraMQ for RocketMQ entrega mensagens ordenadas em sequência estrita FIFO (first-in-first-out). Use mensagens ordenadas quando sua aplicação exigir processamento na ordem exata de envio, como no processamento de pedidos comerciais por horário de submissão ou na sincronização sequencial de alterações em banco de dados.
Este tópico fornece exemplos de código para enviar e receber mensagens ordenadas com o SDK cliente HTTP para Java.
Como funciona
As mensagens ordenadas dividem-se em duas categorias:
Mensagens com ordenação global: Todas as mensagens de um tópico seguem uma única sequência FIFO. Essa abordagem oferece a garantia de ordenação mais forte, mas limita o throughput a uma única partição.
Mensagens com ordenação por partição: As mensagens são distribuídas entre partições com base em uma chave de fragmentação (sharding key). Dentro de cada partição, a ordem FIFO é mantida. Diferentes partições são consumidas de forma independente, o que permite maior throughput sem perder a ordenação onde ela é necessária.
Uma chave de fragmentação determina a qual partição uma mensagem pertence. Mensagens com a mesma chave de fragmentação são sempre entregues à mesma partição e consumidas em ordem. A chave de fragmentação difere da chave da mensagem.
Um broker do ApsaraMQ for RocketMQ define a ordem de geração das mensagens com base na sequência em que o remetente utiliza um único producer ou thread para enviá-las. Caso o remetente use múltiplos producers ou threads para envio simultâneo, a ordem das mensagens será definida pela sequência de recebimento no broker do ApsaraMQ for RocketMQ. Essa ordem pode divergir da sequência de envio original da aplicação.
Ordem de produção e ordem de consumo
A entrega FIFO ponta a ponta depende de dois fatores independentes:
|
Ordem de produção |
Ordem de consumo |
Resultado |
|
Producer único, envios seriais |
Consumo ordenado ( |
Garantia de FIFO por chave de fragmentação |
|
Múltiplos producers ou threads |
Consumo ordenado |
Ordem definida pela chegada ao broker, não pela intenção da aplicação |
|
Producer único, envios seriais |
Consumo concorrente |
Sem garantia de ordem de consumo |
Para garantir ordenação estrita, atenda a ambas as condições: envie a partir de uma única thread de producer e consuma com consumeMessageOrderly.
Pré-requisitos
Antes de começar, certifique-se de ter concluído as seguintes etapas:
Instale o SDK para Java. Para mais informações, consulte Configurar seu ambiente Java.
Crie os recursos a serem especificados no código no console do ApsaraMQ for RocketMQ. Os recursos incluem instâncias, tópicos e grupos de consumidores. Para mais detalhes, consulte Criar recursos.
Obtenha o par de AccessKey da sua conta Alibaba Cloud. Para mais informações, consulte Criar um AccessKey.
Enviar mensagens ordenadas
O exemplo abaixo envia oito mensagens ordenadas distribuídas em duas partições com chaves de fragmentação.
Substitua os placeholders pelos valores reais:
|
Placeholder |
Descrição |
Exemplo |
|
|
Endpoint HTTP obtido na página Instance Details no console do ApsaraMQ for RocketMQ |
|
|
|
Nome do tópico |
|
|
|
ID da instância. Defina como |
|
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 OrderProducer {
public static void main(String[] args) {
MQClient mqClient = new MQClient(
"<your-http-endpoint>",
// Read 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>";
// Set to null or "" if the instance has no namespace
final String instanceId = "<your-instance-id>";
// Get the producer for the specified topic
MQProducer producer;
if (instanceId != null && instanceId != "") {
producer = mqClient.getProducer(instanceId, topic);
} else {
producer = mqClient.getProducer(topic);
}
try {
// Send 8 messages, alternating between 2 sharding keys (partitions)
for (int i = 0; i < 8; i++) {
TopicMessage pubMsg = new TopicMessage(
"hello mq!".getBytes(), // Message body
"A" // Message tag for filtering
);
// Sharding key determines the partition. Messages with the same
// sharding key are delivered to the same partition in FIFO order.
pubMsg.setShardingKey(String.valueOf(i % 2));
pubMsg.getProperties().put("a", String.valueOf(i));
// Publish synchronously. No exception means the send succeeded.
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) {
// Handle send failure: retry or persist the message for later delivery
System.out.println(new Date() + " Send mq message failed. Topic is:" + topic);
e.printStackTrace();
}
mqClient.close();
}
}
Cada tópico suporta apenas um tipo de mensagem. Um tópico criado para mensagens ordenadas não pode enviar ou receber mensagens normais.
Receber mensagens ordenadas
O exemplo a seguir consome mensagens ordenadas com consumeMessageOrderly, método que emprega long polling e garante a ordem FIFO no nível da partição.
Substitua os placeholders pelos valores reais:
|
Placeholder |
Descrição |
Exemplo |
|
|
Endpoint HTTP obtido na página Instance Details |
|
|
|
Nome do tópico |
|
|
|
ID do grupo de consumidores |
|
|
|
ID da instância. Defina como |
|
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 OrderConsumer {
public static void main(String[] args) {
MQClient mqClient = new MQClient(
"<your-http-endpoint>",
System.getenv("ALIBABA_CLOUD_ACCESS_KEY_ID"),
System.getenv("ALIBABA_CLOUD_ACCESS_KEY_SECRET")
);
final String topic = "<your-topic>";
final String groupId = "<your-group-id>";
// Set to null or "" if the instance has no namespace
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);
}
// Continuously poll for messages
do {
List<Message> messages = null;
try {
// consumeMessageOrderly guarantees partition-level FIFO order.
// The consumer must ACK all messages in a batch before the broker
// delivers the next batch from the same partition.
messages = consumer.consumeMessageOrderly(
3, // Max messages per batch (up to 16)
3 // Long polling timeout in seconds (up to 30)
);
} 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
System.out.println("Receive " + messages.size() + " messages:");
for (Message message : messages) {
System.out.println(message);
System.out.println("ShardingKey: " + message.getShardingKey()
+ ", a:" + message.getProperties().get("a"));
}
// ACK consumed messages. If the broker does not receive an ACK before
// the retry timeout, it redelivers the message.
{
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);
}
}
Funcionamento do consumo ordenado
Ao chamar consumeMessageOrderly, o broker impõe os seguintes comportamentos:
Bloqueio no nível da partição: O consumidor pode buscar mensagens de várias partições, mas dentro de cada uma delas as mensagens chegam na ordem de envio.
Controle por confirmação em lote: Confirme (ACK) todas as mensagens do lote atual antes de receber o próximo lote da mesma partição. Se o broker não receber o ACK antes do tempo limite de nova tentativa, ele reenviará a mensagem não confirmada.
Long polling: Quando não há mensagens disponíveis, o broker retém a requisição pelo tempo de polling especificado (até 30 segundos) e responde imediatamente assim que uma mensagem chegar.
Observações de uso
Tópicos de tipo único: Cada tópico aceita somente um tipo de mensagem. Não misture mensagens ordenadas com mensagens normais no mesmo tópico.
Producer único para ordem estrita: Envie todas as mensagens que exigem ordenação relativa a partir de um único producer em uma única thread. O uso de múltiplos producers ou threads concorrentes quebra a garantia de ordenação.
Design da chave de fragmentação: Escolha uma chave que agrupe mensagens relacionadas sem criar partições sobrecarregadas (hot partitions). Por exemplo, utilize um id de pedido para manter todos os eventos de um pedido em sequência, em vez de um id de cliente que poderia concentrar muitas mensagens em uma única partição.
Confirmação rápida: Atrasos na confirmação bloqueiam a entrega de mensagens subsequentes para toda a partição. Processe e confirme (ACK) cada lote o mais rápido possível.
Expiração do receipt handle: Cada receipt handle de mensagem possui um timestamp exclusivo. Se o handle expirar antes da confirmação, o ACK falhará e o broker reenviará a mensagem.
Veja também
Mensagens ordenadas: Garantias de ordenação, tipos de ordenação e conceitos detalhados
Criar recursos: Configure instâncias, tópicos e grupos de consumidores no console