Mensagens ordenadas são entregues e consumidas em ordem estrita FIFO (first-in, first-out). Este tópico fornece códigos de exemplo em Java para enviar e receber mensagens ordenadas via TCP com o Apache RocketMQ SDK for Java.
Funcionamento das mensagens ordenadas
As mensagens ordenadas são particionadas por sharding key. Cada sharding key é mapeada para uma partição dedicada, e as mensagens dentro dessa partição são entregues na exata ordem em que foram enviadas.
Três regras regem esse comportamento:
Ordem dentro da partição -- Mensagens que compartilham a mesma sharding key são sempre entregues em ordem FIFO.
Ausência de ordenação entre partições -- Mensagens com sharding keys diferentes podem chegar em qualquer ordem umas em relação às outras.
Sharding key obrigatória -- Toda mensagem ordenada deve incluir uma sharding key. Essa chave é distinta da chave geral da mensagem usada para consultas.
Utilize sharding keys para modelar seus requisitos de ordenação de negócios. Por exemplo, defina a sharding key como um ID de pedido para que todas as mensagens do mesmo pedido sejam processadas sequencialmente, enquanto mensagens de pedidos diferentes são processadas em paralelo.
Para obter mais informações sobre semântica e restrições de ordenação, consulte Mensagens ordenadas.
Pré-requisitos
Antes de começar, verifique se você possui:
O SDK for Java 1.2.7 ou posterior. Consulte Notas de versão
Um ambiente preparado. Consulte Preparar o ambiente
(Opcional) Configurações de log definidas. Consulte Configurações de log
Se você está começando agora com o ApsaraMQ for RocketMQ, inicie pelo projeto de demonstração para configurar um projeto funcional antes de enviar e receber mensagens.
Enviar mensagens ordenadas
Garantia de ordenação
O broker determina a ordem das mensagens com base na sequência em que o remetente usa um único producer ou thread para enviá-las. Se múltiplos producers ou threads enviarem mensagens simultaneamente, o broker as ordenará pelo horário de chegada, o que pode divergir da ordem de negócios pretendida.
Código de exemplo
Substitua os placeholders abaixo pelos seus valores reais:
|
Placeholder |
Descrição |
Exemplo |
|
|
Group ID criado no console do ApsaraMQ for RocketMQ |
GID_order_group |
|
|
AccessKey ID para autenticação |
LTAI5tXxx |
|
|
AccessKey secret para autenticação |
xXxXxXx |
|
|
Endpoint TCP obtido na seção TCP Endpoint da página Instance Details |
http://MQ_INST_xxx.mq-internet-access.mq-internet.aliyuncs.com:80 |
Para acessar o código completo, consulte o repositório de código do ApsaraMQ for RocketMQ.
package com.aliyun.openservices.ons.example.order;
import com.aliyun.openservices.ons.api.Message;
import com.aliyun.openservices.ons.api.ONSFactory;
import com.aliyun.openservices.ons.api.PropertyKeyConst;
import com.aliyun.openservices.ons.api.SendResult;
import com.aliyun.openservices.ons.api.order.OrderProducer;
import java.util.Date;
import java.util.Properties;
public class ProducerClient {
public static void main(String[] args) {
Properties properties = new Properties();
// Group ID created in the ApsaraMQ for RocketMQ console.
properties.put(PropertyKeyConst.GROUP_ID, "<your-group-id>");
// AccessKey ID for authentication.
properties.put(PropertyKeyConst.AccessKey, "<your-access-key>");
// AccessKey secret for authentication.
properties.put(PropertyKeyConst.SecretKey, "<your-secret-key>");
// TCP endpoint from the Instance Details page in the ApsaraMQ for RocketMQ console.
properties.put(PropertyKeyConst.NAMESRV_ADDR, "<your-tcp-endpoint>");
OrderProducer producer = ONSFactory.createOrderProducer(properties);
// Call start() once before sending any messages.
producer.start();
for (int i = 0; i < 1000; i++) {
String orderId = "biz_" + i % 10;
Message msg = new Message(
// Topic for ordered messages.
"Order_global_topic",
// Tag for filtering on the broker side.
"TagA",
// Message body in binary format.
// Producer and consumer must agree on serialization and deserialization.
"send order global msg".getBytes()
);
// Message key -- a business identifier, ideally globally unique.
// Use this key to look up the message in the ApsaraMQ for RocketMQ console.
msg.setKey(orderId);
// Sharding key -- determines the partition for ordered delivery.
// Messages with the same sharding key are delivered in FIFO order.
String shardingKey = String.valueOf(orderId);
try {
SendResult sendResult = producer.send(msg, shardingKey);
if (sendResult != null) {
System.out.println(new Date() + " Send mq message success. Topic is:" + msg.getTopic() + " msgId is: " + sendResult.getMessageId());
}
} catch (Exception e) {
// Handle failure: retry or persist the message for later resend.
System.out.println(new Date() + " Send mq message failed. Topic is:" + msg.getTopic());
e.printStackTrace();
}
}
// Shut down the producer before exiting.
// Skip this if the producer sends messages frequently.
producer.shutdown();
}
}
Receber mensagens ordenadas
Comportamento de nova tentativa
Quando o consumo de uma mensagem ordenada falha, o broker bloqueia todas as mensagens subsequentes na mesma partição até que a mensagem atual seja processada com sucesso ou atinja o número máximo de tentativas. Esse bloqueio preserva a ordenação FIFO, mas pode atrasar o processamento de outras mensagens na partição.
Dois parâmetros controlam o comportamento de nova tentativa:
|
Parâmetro |
Descrição |
Valores válidos |
|
|
Tempo de espera em milissegundos entre tentativas para uma mensagem ordenada com falha |
10 -- 30.000 |
|
|
Número máximo de tentativas antes de a mensagem ser ignorada |
Inteiro (por exemplo, |
Defina esses valores conforme sua tolerância a atrasos no processamento. Um valor menor para SuspendTimeMillis acelera as novas tentativas, porém aumenta a carga no broker. Já um valor menor para MaxReconsumeTimes desbloqueia a partição mais rapidamente, mas pode ignorar mensagens que teriam sucesso com mais tentativas.
Código de exemplo
O código de exemplo a seguir utiliza os mesmos placeholders do exemplo do producer. Substitua-os pelos seus valores reais antes de executar o código.
package com.aliyun.openservices.ons.example.order;
import com.aliyun.openservices.ons.api.Message;
import com.aliyun.openservices.ons.api.ONSFactory;
import com.aliyun.openservices.ons.api.PropertyKeyConst;
import com.aliyun.openservices.ons.api.order.ConsumeOrderContext;
import com.aliyun.openservices.ons.api.order.MessageOrderListener;
import com.aliyun.openservices.ons.api.order.OrderAction;
import com.aliyun.openservices.ons.api.order.OrderConsumer;
import java.util.Properties;
public class ConsumerClient {
public static void main(String[] args) {
Properties properties = new Properties();
// Consumer group ID created in the ApsaraMQ for RocketMQ console.
properties.put(PropertyKeyConst.GROUP_ID, "<your-group-id>");
// AccessKey ID for authentication.
properties.put(PropertyKeyConst.AccessKey, "<your-access-key>");
// AccessKey secret for authentication.
properties.put(PropertyKeyConst.SecretKey, "<your-secret-key>");
// TCP endpoint from the Instance Details page in the ApsaraMQ for RocketMQ console.
properties.put(PropertyKeyConst.NAMESRV_ADDR, "<your-tcp-endpoint>");
// Wait time (ms) before retrying a failed ordered message. Valid values: 10 to 30,000.
properties.put(PropertyKeyConst.SuspendTimeMillis, "100");
// Maximum retry attempts for a failed message.
properties.put(PropertyKeyConst.MaxReconsumeTimes, "20");
OrderConsumer consumer = ONSFactory.createOrderedConsumer(properties);
consumer.subscribe(
// Topic to subscribe to.
"Order_global_topic",
// Tag filter expression:
// "*" -- subscribe to all tags
// "TagA || TagB || TagC" -- subscribe to TagA, TagB, or TagC only
"*",
new MessageOrderListener() {
/**
* Return OrderAction.Success after processing the message.
* Return OrderAction.Suspend if processing fails or throws
* an exception -- the broker will retry after SuspendTimeMillis.
*/
@Override
public OrderAction consume(Message message, ConsumeOrderContext context) {
System.out.println(message);
return OrderAction.Success;
}
});
// Call start() once to begin consuming.
consumer.start();
}
}
Próximos passos
Conheça os conceitos e restrições das mensagens ordenadas.
Explore o repositório de código completo do ApsaraMQ for RocketMQ para encontrar exemplos adicionais.