As mensagens ordenadas no ApsaraMQ for RocketMQ garantem entrega estrita na ordem FIFO (first-in, first-out). Use esse tipo de mensagem quando a sequência de processamento for crítica para sua lógica de negócios — por exemplo, em correspondência de transações em que a primeira oferta em um determinado preço deve prevalecer, ou na sincronização de alterações de banco de dados em que inserções, atualizações e exclusões precisam ser replicadas na ordem exata.
Este tópico demonstra como enviar e receber mensagens ordenadas com o pacote @aliyunmq/mq-http-sdk para Node.js.
Funcionamento da ordenação
O ApsaraMQ for RocketMQ oferece dois escopos de ordenação:
Mensagens globalmente ordenadas — Todas as mensagens de um tópico são entregues e consumidas em uma única sequência FIFO.
Mensagens ordenadas 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 entrega e o consumo seguem a ordem FIFO. Mensagens em partições diferentes podem ser consumidas simultaneamente.
A chave de fragmentação define para qual partição a mensagem será roteada. Já a chave de mensagem é um índice separado, usado para localização. Não confunda esses dois conceitos.
A garantia de ordenação depende de dois fatores independentes:
|
Garantia |
Requisito |
|
Ordem de produção |
Envie as mensagens a partir de um único producer em uma única thread. Caso múltiplos producers ou threads enviem dados simultaneamente, o broker registrará as mensagens pela ordem de chegada, o que pode divergir da ordem lógica de negócios desejada. |
|
Ordem de consumo |
Confirme (ACK) cada lote de mensagens antes de consumir o próximo lote da mesma partição. Se o broker não receber a confirmação, ele reentregará a mensagem pendente antes de liberar novas mensagens. |
Pré-requisitos
Antes de começar, verifique se você já:
Crie os recursos no console do ApsaraMQ for RocketMQ: uma instância, um tópico e um grupo de consumidores
Crie um par de AccessKey para sua conta Alibaba Cloud
Enviar mensagens ordenadas
A ordenação só é garantida quando um único producer envia mensagens em uma única thread. Se múltiplos producers ou threads enviarem dados ao mesmo tempo, o broker registrará as mensagens conforme a ordem de chegada, podendo alterar a sequência lógica esperada pelo negócio.
Substitua os placeholders abaixo pelos valores reais antes de executar o código:
|
Placeholder |
Descrição |
|
|
Endpoint HTTP obtido na página Instance Details > seção HTTP Endpoint no console do ApsaraMQ for RocketMQ |
|
|
Nome do tópico criado no console |
|
|
ID da instância. Se a instância possuir namespace, especifique o ID. Caso contrário, defina como |
const {
MQClient,
MessageProperties
} = require('@aliyunmq/mq-http-sdk');
// HTTP endpoint. Get this from Instance Details > HTTP Endpoint in the console.
const endpoint = "${HTTP_ENDPOINT}";
// Load credentials from environment variables.
const accessKeyId = process.env['ALIBABA_CLOUD_ACCESS_KEY_ID'];
const accessKeySecret = process.env['ALIBABA_CLOUD_ACCESS_KEY_SECRET'];
var client = new MQClient(endpoint, accessKeyId, accessKeySecret);
// Topic created in the ApsaraMQ for RocketMQ console.
const topic = "${TOPIC}";
// Instance ID. Set to null or "" if the instance does not have a namespace.
const instanceId = "${INSTANCE_ID}";
const producer = client.getProducer(instanceId, topic);
(async function(){
try {
// Send 8 messages. Sharding key i % 2 routes messages to 2 partitions.
for(var i = 0; i < 8; i++) {
msgProps = new MessageProperties();
// Set a custom property.
msgProps.putProperty("a", i);
// Set the sharding key. Messages with the same key go to the same partition.
msgProps.shardingKey(i % 2);
// Publish with body "hello mq." and tag "TagA".
res = await producer.publishMessage("hello mq.", "TagA", msgProps);
console.log("Publish message: MessageID:%s,BodyMD5:%s", res.body.MessageId, res.body.MessageBodyMD5);
}
} catch(e) {
// Handle send failures. Implement retry or persistence logic as needed.
console.log(e)
}
})();
Receber mensagens ordenadas
Para consumir mensagens mantendo a ordem, use o método consumeMessageOrderly() em vez de consumeMessage(). Essa abordagem assegura que as mensagens da mesma partição cheguem na sequência em que foram enviadas.
Comportamentos principais:
Uma única chamada pode retornar mensagens de várias partições, mas a ordem interna de cada partição é sempre preservada.
Caso o broker não receba um ACK, ele reentrega a mensagem pendente antes de liberar a próxima mensagem daquela partição.
Confirme todas as mensagens de um lote antes que o broker entregue o próximo lote da mesma partição.
Long polling: se não houver mensagens disponíveis, o broker mantém a conexão aberta durante o tempo de espera especificado e responde imediatamente assim que uma nova mensagem chegar.
const {
MQClient,
} = require('@aliyunmq/mq-http-sdk');
// HTTP endpoint. Get this from Instance Details > HTTP Endpoint in the console.
const endpoint = "${HTTP_ENDPOINT}";
// Load credentials from environment variables.
const accessKeyId = process.env['ALIBABA_CLOUD_ACCESS_KEY_ID'];
const accessKeySecret = process.env['ALIBABA_CLOUD_ACCESS_KEY_SECRET'];
var client = new MQClient(endpoint, accessKeyId, accessKeySecret);
// Topic created in the ApsaraMQ for RocketMQ console.
const topic = "${TOPIC}";
// Consumer group ID created in the ApsaraMQ for RocketMQ console.
const groupId = "GID_http";
// Instance ID. Set to null or "" if the instance does not have a namespace.
const instanceId = "${INSTANCE_ID}";
const consumer = client.getConsumer(instanceId, topic, groupId);
(async function(){
while(true) {
try {
// consumeMessageOrderly(numOfMessages, waitSeconds)
// numOfMessages: max messages per call (up to 16)
// waitSeconds: long polling duration in seconds (up to 30)
res = await consumer.consumeMessageOrderly(
3, // Receive up to 3 messages per call.
3 // Wait up to 3 seconds if no messages are available.
);
if (res.code == 200) {
console.log("Consume Messages, requestId:%s", res.requestId);
const handles = res.body.map((message) => {
console.log("\tMessageId:%s,Tag:%s,PublishTime:%d,NextConsumeTime:%d,FirstConsumeTime:%d,ConsumedTimes:%d,Body:%s" +
",Props:%j,ShardingKey:%s,Prop-A:%s,Tag:%s",
message.MessageId, message.MessageTag, message.PublishTime, message.NextConsumeTime, message.FirstConsumeTime, message.ConsumedTimes,
message.MessageBody,message.Properties,message.ShardingKey,message.Properties.a);
return message.ReceiptHandle;
});
// Acknowledge consumed messages.
// If NextConsumeTime passes without an ACK, the broker redelivers the message.
res = await consumer.ackMessage(handles);
if (res.code != 204) {
// ACK failed -- typically caused by an expired receipt handle.
console.log("Ack Message Fail:");
const failHandles = res.body.map((error)=>{
console.log("\tErrorHandle:%s, Code:%s, Reason:%s\n", error.ReceiptHandle, error.ErrorCode, error.ErrorMessage);
return error.ReceiptHandle;
});
handles.forEach((handle)=>{
if (failHandles.indexOf(handle) < 0) {
console.log("\tSucHandle:%s\n", handle);
}
});
} else {
console.log("Ack Message suc, RequestId:%s\n\t", res.requestId, handles.join(','));
}
}
} catch(e) {
if (e.Code.indexOf("MessageNotExist") > -1) {
// No messages available. Long polling continues.
console.log("Consume Message: no new message, RequestId:%s, Code:%s", e.RequestId, e.Code);
} else {
console.log(e);
}
}
}
})();
Melhores práticas
Use uma única thread de producer para envio. O uso de múltiplos producers ou threads para a mesma chave de fragmentação pode quebrar a ordem esperada, pois o broker registra as mensagens pelo horário de chegada, e não pelo horário de envio.
Confirme todos os lotes antes de prosseguir. O broker bloqueia a entrega do próximo lote de uma partição até que todas as mensagens do lote atual sejam confirmadas. Ignorar ou atrasar o envio de ACKs paralisa o consumo.