Todos os produtos
Search
Central de documentação

ApsaraMQ for RocketMQ:Send and receive ordered messages

Última atualização: Jun 27, 2026

As mensagens ordenadas no ApsaraMQ for RocketMQ garantem entrega estrita do tipo first-in-first-out (FIFO) dentro de um tópico ou partição. Utilize mensagens ordenadas quando a sequência de processamento afetar sua lógica de negócios, como na correspondência de transações em que o primeiro lance em um determinado preço deve vencer, ou na sincronização de alterações de banco de dados em que as operações de inserção, atualização e exclusão devem ser reproduzidas na ordem correta.

Os exemplos a seguir em C++ demonstram como enviar e receber mensagens ordenadas por meio do SDK cliente HTTP.

Funcionamento da ordenação

O ApsaraMQ for RocketMQ oferece suporte a dois escopos de ordenação:

  • Mensagens globalmente ordenadas -- Todas as mensagens de um tópico são entregues em ordem FIFO. Adote este escopo quando cada mensagem precisar ser processada sequencialmente.

  • Mensagens ordenadas por partição -- As mensagens são distribuídas entre partições com base na chave de fragmentação (sharding key). Dentro de cada partição, a entrega segue a ordem FIFO. Este escopo é ideal para obter maior throughput quando apenas mensagens relacionadas exigem ordenação, como agrupamentos por id de pedido ou id de usuário.

A chave de fragmentação determina qual partição recebe uma mensagem. Mensagens com a mesma chave de fragmentação sempre vão para a mesma partição e são consumidas na ordem de envio. Note que a chave de fragmentação difere da chave da mensagem.

Para mais informações, consulte Mensagens ordenadas.

Ordem de produção

O broker define a ordem das mensagens com base na sequência em que um único produtor ou thread as envia. Para preservar a ordenação:

  • Envie mensagens relacionadas a partir de um único produtor em uma única thread.

  • Caso múltiplos produtores ou threads enviem mensagens simultaneamente, o broker as ordenará pelo horário de chegada, o que pode divergir da sequência de negócios pretendida.

Ordem de consumo

O consumidor busca mensagens ordenadas por partição em uma ou mais partições e processa as mensagens de cada partição na ordem de envio. A garantia de ordenação depende do reconhecimento adequado:

  • O consumidor precisa confirmar todas as mensagens de um lote antes de buscar o próximo lote na mesma partição.

  • Se o broker não receber um reconhecimento (ACK) de uma mensagem antes do tempo limite, ele reenviará essa mensagem.

  • Cada reentrega gera um novo identificador de recibo (receipt handle) com um timestamp exclusivo.

Pré-requisitos

Antes de começar, verifique se você possui:

Enviar mensagens ordenadas

Importante

O broker sequencia as mensagens pela ordem de chegada de um único produtor ou thread. Se múltiplos produtores ou threads enviarem mensagens simultaneamente, a ordem no lado do broker poderá diferir da ordem de envio definida no lado da aplicação.

Substitua os placeholders a seguir pelos seus valores reais:

Placeholder

Descrição

Onde encontrar

${HTTP_ENDPOINT}

Endpoint HTTP da sua instância

Página Instance Details > HTTP Endpoint no console do ApsaraMQ for RocketMQ

${TOPIC}

Nome do tópico

O tópico que você criou no console

${INSTANCE_ID}

id da instância

Página Instance Details no console. Defina como uma string vazia se a instância não tiver namespace

//#include <iostream>
#include <fstream>
#include <time.h>
#include "mq_http_sdk/mq_client.h"

using namespace std;
using namespace mq::http::sdk;

int main() {

    MQClient mqClient(
            // The HTTP endpoint. You can obtain the endpoint in the HTTP Endpoint section of the Instance Details page in the ApsaraMQ for RocketMQ console.
            "${HTTP_ENDPOINT}",
            // Make sure that the environment variables ALIBABA_CLOUD_ACCESS_KEY_ID and ALIBABA_CLOUD_ACCESS_KEY_SECRET are configured.
	          // The AccessKey ID that is used for authentication.
	          System.getenv("ALIBABA_CLOUD_ACCESS_KEY_ID"),
	          // The AccessKey secret that is used for authentication.
	          System.getenv("ALIBABA_CLOUD_ACCESS_KEY_SECRET")
            );

    // The topic in which the message is produced. You must create the topic in the ApsaraMQ for RocketMQ console.
    string topic = "${TOPIC}";
    // The ID of the instance to which the topic belongs. You must create the instance in the ApsaraMQ for RocketMQ console.
    // If the instance has a namespace, specify the ID of the instance. If the instance does not have a namespace, set the instanceID parameter to null or an empty string. You can obtain the namespace of the instance on the Instance Details page in the ApsaraMQ for RocketMQ console.
    string instanceId = "${INSTANCE_ID}";

    MQProducerPtr producer;
    if (instanceId == "") {
        producer = mqClient.getProducerRef(topic);
    } else {
        producer = mqClient.getProducerRef(instanceId, topic);
    }

    try {
        // Cyclically send four messages.
        for (int i = 0; i < 8; i++)
        {
            PublishMessageResponse pmResp;
            // The message content.
            TopicMessage pubMsg("Hello, mq!order msg!");
            // The sharding key that is used to distribute ordered messages to a specific partition. Sharding keys can be used to identify partitions. A sharding key is different from a message key.
            pubMsg.setShardingKey(std::to_string(i % 2));
            // The custom attributes of the message.
            pubMsg.putProperty("a",std::to_string(i));
            producer->publishMessage(pubMsg, pmResp);
            cout << "Publish mq message success. Topic is: " << topic
                << ", msgId is:" << pmResp.getMessageId()
                << ", bodyMD5 is:" << pmResp.getMessageBodyMD5() << endl;
        }
    } catch (MQServerException& me) {
        cout << "Request Failed: " + me.GetErrorCode() << ", requestId is:" << me.GetRequestId() << endl;
        return -1;
    } catch (MQExceptionBase& mb) {
        cout << "Request Failed: " + mb.ToString() << endl;
        return -2;
    }

    return 0;
}

Este exemplo envia oito mensagens usando duas chaves de fragmentação (0 e 1 via i % 2). As mensagens com a mesma chave de fragmentação são entregues à mesma partição na ordem de envio. A chamada putProperty anexa um atributo personalizado a cada mensagem para fins de rastreamento.

Receber mensagens ordenadas

Invoque consumeMessageOrderly em vez do método de consumo padrão para manter a ordenação no nível da partição. No modo de long polling, o broker mantém a solicitação aberta durante a duração especificada caso nenhuma mensagem esteja disponível e responde imediatamente assim que uma mensagem chega.

Além dos placeholders listados acima, substitua o seguinte:

Placeholder

Descrição

Onde encontrar

${GROUP_ID}

id do grupo de consumidores

O grupo de consumidores criado no console do ApsaraMQ for RocketMQ

#include <vector>
#include <fstream>
#include "mq_http_sdk/mq_client.h"

#ifdef _WIN32
#include <windows.h>
#else
#include <unistd.h>
#endif

using namespace std;
using namespace mq::http::sdk;

int main() {

    MQClient mqClient(
            // The HTTP endpoint. You can obtain the endpoint in the HTTP Endpoint section of the Instance Details page in the ApsaraMQ for RocketMQ console.
            "${HTTP_ENDPOINT}",
            // Make sure that the environment variables ALIBABA_CLOUD_ACCESS_KEY_ID and ALIBABA_CLOUD_ACCESS_KEY_SECRET are configured.
	          // The AccessKey ID that is used for authentication.
	          System.getenv("ALIBABA_CLOUD_ACCESS_KEY_ID"),
	          // The AccessKey secret that is used for authentication.
	          System.getenv("ALIBABA_CLOUD_ACCESS_KEY_SECRET")
            );

    // The topic in which the message is produced. You must create the topic in the ApsaraMQ for RocketMQ console.
    string topic = "${TOPIC}";
    // The ID of the consumer group that you created in the ApsaraMQ for RocketMQ console.
    string groupId = "${GROUP_ID}";
    // The ID of the instance to which the topic belongs. You must create the instance in the ApsaraMQ for RocketMQ console.
    // If the instance has a namespace, specify the ID of the instance. If the instance does not have a namespace, set the instanceID parameter to null or an empty string. You can obtain the namespace of the instance on the Instance Details page in the ApsaraMQ for RocketMQ console.
    string instanceId = "${INSTANCE_ID}";

    MQConsumerPtr consumer;
    if (instanceId == "") {
        consumer = mqClient.getConsumerRef(topic, groupId);
    } else {
        consumer = mqClient.getConsumerRef(instanceId, topic, groupId, "");
    }

    do {
        try {
            std::vector<Message> messages;
            // Consume messages in long polling mode. The consumer may pull partitionally ordered messages from multiple partitions. The consumer consumes messages from the same partition in the order in which the messages are sent.
            // Assume that a consumer pulls partitionally ordered messages from one partition. If the broker fails to receive an acknowledgement (ACK) for a message from the consumer, the broker delivers the message in the partition to the consumer again.
            // The consumer can consume the next batch of messages from a partition only after all messages that are pulled from the partition in the previous batch are acknowledged as consumed.
            // In long polling mode, if no message in the topic is available for consumption, the request is suspended on the broker for the specified period of time. If a message becomes available for consumption within the specified period of time, the broker immediately sends a response to the consumer. In this example, the value is specified as 3 seconds.
            consumer->consumeMessageOrderly(
                    3, // The maximum number of messages that can be consumed at a time. In this example, the value is specified as 3. The maximum value that you can specify is 16.
                    3, // The duration of a long polling cycle. Unit: seconds. In this example, the value is specified as 3. The maximum value that you can specify is 30.
                    messages
            );
            cout << "Consume: " << messages.size() << " Messages!" << endl;

            // The message consumption logic.
            std::vector<std::string> receiptHandles;
            for (std::vector<Message>::iterator iter = messages.begin();
                    iter != messages.end(); ++iter)
            {
                cout << "MessageId: " << iter->getMessageId()
                    << " PublishTime: " << iter->getPublishTime()
                    << " Tag: " << iter->getMessageTag()
                    << " Body: " << iter->getMessageBody()
                    << " FirstConsumeTime: " << iter->getFirstConsumeTime()
                    << " NextConsumeTime: " << iter->getNextConsumeTime()
                    << " ConsumedTimes: " << iter->getConsumedTimes()
                    << " Properties: " << iter->getPropertiesAsString()
                    << " ShardingKey: " << iter->getShardingKey() << endl;
                receiptHandles.push_back(iter->getReceiptHandle());
            }

            // Obtain an ACK from the consumer.
            // If the broker fails to receive an ACK for a message from the consumer before the period of time that is specified by the Message.NextConsumeTime parameter elapses, the broker delivers the message for consumption again.
            // A unique timestamp is specified for the handle of a message each time the message is consumed.
            AckMessageResponse bdmResp;
            consumer->ackMessage(receiptHandles, bdmResp);
            if (!bdmResp.isSuccess()) {
                // If the handle of a message times out, the broker cannot receive an ACK for the message from the consumer.
                const std::vector<AckMessageFailedItem>& failedItems =
                    bdmResp.getAckMessageFailedItem();
                for (std::vector<AckMessageFailedItem>::const_iterator iter = failedItems.begin();
                        iter != failedItems.end(); ++iter)
                {
                    cout << "AckFailedItem: " << iter->errorCode
                        << "  " << iter->receiptHandle << endl;
                }
            } else {
                cout << "Ack: " << messages.size() << " messages suc!" << endl;
            }
        } catch (MQServerException& me) {
            if (me.GetErrorCode() == "MessageNotExist") {
                cout << "No message to consume! RequestId: " + me.GetRequestId() << endl;
                continue;
            }
            cout << "Request Failed: " + me.GetErrorCode() + ".RequestId: " + me.GetRequestId() << endl;
#ifdef _WIN32
            Sleep(2000);
#else
            usleep(2000 * 1000);
#endif
        } catch (MQExceptionBase& mb) {
            cout << "Request Failed: " + mb.ToString() << endl;
#ifdef _WIN32
            Sleep(2000);
#else
            usleep(2000 * 1000);
#endif
        }

    } while(true);
}

Parâmetros de consumeMessageOrderly

Parâmetro

Descrição

Intervalo

Tamanho do lote

Número máximo de mensagens por busca

Máximo: 16

Duração do long polling

Segundos em que o broker retém a solicitação quando não há mensagens disponíveis

Máximo: 30

Comportamento de reconhecimento

Após processar cada lote, chame ackMessage com os identificadores de recibo de todas as mensagens processadas.

  • Bloqueio por lote: O consumidor não consegue buscar o próximo lote de uma partição até que todas as mensagens do lote atual sejam reconhecidas.

  • Reentrega por timeout: Caso o broker não receba um ACK antes de NextConsumeTime, ele reenvia a mensagem com um novo identificador de recibo contendo um timestamp exclusivo.

  • **Tratamento de MessageNotExist:** Este código de erro indica que não há mensagens disponíveis. O exemplo lida com isso continuando o loop de polling sem atraso.

  • Recuperação de erros: Em outras exceções, o exemplo aguarda 2 segundos antes de tentar novamente para evitar loops de erro contínuos.

Observações de uso

  • Um único produtor por escopo de ordenação. Envie todas as mensagens que exigem ordenação relativa a partir de um único produtor em uma única thread. Envios com múltiplos produtores ou threads podem quebrar a ordem pretendida.

  • Projete chaves de fragmentação com a granularidade adequada. Utilize chaves significativas para o negócio, como ids de pedidos ou ids de usuários. Evite concentrar muitas mensagens sob uma única chave de fragmentação, pois isso sobrecarrega uma partição e limita o throughput.

  • Faça o reconhecimento prontamente. O atraso no reconhecimento bloqueia o consumo das mensagens subsequentes na mesma partição, aumentando a latência de ponta a ponta.

  • Lide com falhas sem quebrar a ordem. Se o processamento da mensagem falhar, permita que a mensagem atinja o timeout e seja reentregue, em vez de confirmá-la e pular adiante, o que violaria a garantia de ordenação.