Todos os produtos
Search
Central de documentação

ApsaraMQ for RocketMQ:Send and receive ordered messages

Última atualização: Jun 27, 2026

O ApsaraMQ for RocketMQ entrega mensagens ordenadas seguindo rigorosamente a ordem FIFO (first-in-first-out). Este tópico fornece códigos de exemplo em PHP para enviar e receber mensagens ordenadas por meio do SDK cliente HTTP.

Tipos de ordenação

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

Tipo

Escopo de ordenação

Funcionamento

Ordenação global

Todo o tópico

Todas as mensagens do tópico são enviadas e consumidas em ordem FIFO

Ordenação por partição

Por partição

As mensagens são distribuídas nas partições com base na chave de sharding. Dentro de cada partição, o consumo ocorre em ordem FIFO

A chave de sharding identifica a qual partição uma mensagem pertence. Ela difere da chave da mensagem.

Para mais informações, consulte Mensagens ordenadas.

Garantias de ordenação

A entrega ordenada depende do comportamento correto tanto do produtor quanto do consumidor.

Importante

O broker determina a ordem das mensagens com base na sequência em que um único produtor ou thread as envia. Se vários produtores ou threads enviarem mensagens simultaneamente, a ordem será definida pela chegada ao broker, o que pode divergir da ordem lógica de negócios pretendida.

Lado do consumidor:

  • O consumidor pode buscar mensagens ordenadas de várias partições em um único lote. As mensagens de cada partição mantêm a ordem de envio original.

  • Confirme (ACK) todas as mensagens de uma partição antes de receber o próximo lote dessa mesma partição.

  • Caso o broker não receba um ACK antes do tempo limite definido por getNextConsumeTime(), a mensagem será reentregue.

Pré-requisitos

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

Enviar mensagens ordenadas

O código abaixo envia quatro mensagens ordenadas e as distribui em duas partições usando setShardingKey(). Mensagens que compartilham a mesma chave de sharding são entregues em ordem FIFO dentro dessa partição.

<?php

require "vendor/autoload.php";

use MQ\Model\TopicMessage;
use MQ\MQClient;

class ProducerTest
{
    private $client;
    private $producer;

    public function __construct()
    {
        $this->client = new MQClient(
            // HTTP endpoint. Find this in the HTTP Endpoint section on the Instance Details
            // page in the ApsaraMQ for RocketMQ console.
            "${HTTP_ENDPOINT}",
            // Get AccessKey credentials from environment variables.
            getenv("ALIBABA_CLOUD_ACCESS_KEY_ID"),
            getenv('ALIBABA_CLOUD_ACCESS_KEY_SECRET')
        );

        // Topic created in the ApsaraMQ for RocketMQ console.
        $topic = "${TOPIC}";
        // Instance ID. If the instance has no namespace, set this to null or "".
        // Check the Instance Details page for namespace information.
        $instanceId = "${INSTANCE_ID}";

        $this->producer = $this->client->getProducer($instanceId, $topic);
    }

    public function run()
    {
        try
        {
            for ($i = 1; $i <= 4; $i++)
            {
                $publishMessage = new TopicMessage(
                    "hello mq! " // Message body
                );
                // Set a custom message property.
                $publishMessage->putProperty("a", $i);
                // Set the sharding key to distribute messages across partitions.
                // Messages with the same sharding key are delivered in FIFO order.
                $publishMessage->setShardingKey($i % 2);

                $result = $this->producer->publishMessage($publishMessage);

                print "Send mq message success. msgId is:" . $result->getMessageId() . ", bodyMD5 is:" . $result->getMessageBodyMD5() . "\n";
            }
        } catch (\Exception $e) {
            print_r($e->getMessage() . "\n");
        }
    }
}

$instance = new ProducerTest();
$instance->run();

?>

Substitua os placeholders a seguir pelos valores reais:

Placeholder

Descrição

${HTTP_ENDPOINT}

Endpoint HTTP obtido na página de detalhes da instância

${TOPIC}

Nome do tópico

${INSTANCE_ID}

ID da instância (null ou "" se não houver namespace)

Receber mensagens ordenadas

O código a seguir consome mensagens ordenadas utilizando consumeMessageOrderly(), preservando a ordem FIFO dentro de cada partição. O consumidor emprega long polling: se nenhuma mensagem estiver disponível, a solicitação permanece no broker até a chegada de uma nova mensagem ou o esgotamento do tempo limite de polling.

<?php

require "vendor/autoload.php";

use MQ\MQClient;

class ConsumerTest
{
    private $client;
    private $consumer;

    public function __construct()
    {
        $this->client = new MQClient(
            // HTTP endpoint. Find this in the HTTP Endpoint section on the Instance Details
            // page in the ApsaraMQ for RocketMQ console.
            "${HTTP_ENDPOINT}",
            // Get AccessKey credentials from environment variables.
            getenv("ALIBABA_CLOUD_ACCESS_KEY_ID"),
            getenv('ALIBABA_CLOUD_ACCESS_KEY_SECRET')
        );

        // Topic created in the ApsaraMQ for RocketMQ console.
        $topic = "${TOPIC}";
        // Consumer group ID created in the ApsaraMQ for RocketMQ console.
        $groupId = "${GROUP_ID}";
        // Instance ID. If the instance has no namespace, set this to null or "".
        // Check the Instance Details page for namespace information.
        $instanceId = "${INSTANCE_ID}";

        $this->consumer = $this->client->getConsumer($instanceId, $topic, $groupId);
    }

    public function ackMessages($receiptHandles)
    {
        try {
            $this->consumer->ackMessage($receiptHandles);
        } catch (\Exception $e) {
            if ($e instanceof MQ\Exception\AckMessageException) {
                // Handle ACK failures. This occurs when a receipt handle has expired.
                printf("Ack Error, RequestId:%s\n", $e->getRequestId());
                foreach ($e->getAckMessageErrorItems() as $errorItem) {
                    printf("\tReceiptHandle:%s, ErrorCode:%s, ErrorMsg:%s\n", $errorItem->getReceiptHandle(), $errorItem->getErrorCode(), $errorItem->getErrorCode());
                }
            }
        }
    }

    public function run()
    {
        // Poll for messages continuously. For production use, run multiple threads
        // to consume messages concurrently across partitions.
        while (True) {
            try {
                // consumeMessageOrderly() ensures FIFO order within each partition.
                // If the broker does not receive an ACK for a message, it redelivers
                // that message before sending subsequent messages from the same partition.
                $messages = $this->consumer->consumeMessageOrderly(
                    3, // Max messages per batch (up to 16)
                    3  // Long polling timeout in seconds (up to 30)
                );
            } catch (\MQ\Exception\MessageResolveException $e) {
                // Handle messages that cannot be parsed due to invalid characters.
                $messages = $e->getPartialResult()->getMessages();
                $failMessages = $e->getPartialResult()->getFailResolveMessages();

                $receiptHandles = array();
                foreach ($messages as $message) {
                    $receiptHandles[] = $message->getReceiptHandle();
                    printf("MsgID %s\n", $message->getMessageId());
                }
                foreach ($failMessages as $failMessage) {
                    $receiptHandles[] = $failMessage->getReceiptHandle();
                    printf("Fail To Resolve Message. MsgID %s\n", $failMessage->getMessageId());
                }
                $this->ackMessages($receiptHandles);
                continue;
            } catch (\Exception $e) {
                if ($e instanceof MQ\Exception\MessageNotExistException) {
                    // No messages available. Long polling continues automatically.
                    printf("No message, contine long polling!RequestId:%s\n", $e->getRequestId());
                    continue;
                }

                print_r($e->getMessage() . "\n");

                sleep(3);
                continue;
            }

            print "======>consume finish, messages:\n";

            $receiptHandles = array();
            foreach ($messages as $message) {
                $receiptHandles[] = $message->getReceiptHandle();
                printf("MessageID:%s TAG:%s BODY:%s \nPublishTime:%d, FirstConsumeTime:%d, \nConsumedTimes:%d, NextConsumeTime:%d,ShardingKey:%s\n",
                    $message->getMessageId(), $message->getMessageTag(), $message->getMessageBody(),
                    $message->getPublishTime(), $message->getFirstConsumeTime(), $message->getConsumedTimes(), $message->getNextConsumeTime(),
                    $message->getShardingKey());
                print_r($message->getProperties());
            }

            // ACK all messages. If the broker does not receive an ACK before
            // getNextConsumeTime(), it redelivers the message. Each redelivery
            // assigns a new receipt handle with an updated timestamp.
            print_r($receiptHandles);
            $this->ackMessages($receiptHandles);
            print "=======>ack finish\n";

        }

    }
}

$instance = new ConsumerTest();
$instance->run();

?>

Além dos placeholders mencionados na seção Enviar mensagens ordenadas, substitua o seguinte:

Placeholder

Descrição

${GROUP_ID}

ID do grupo de consumidores

Parâmetros de consumeMessageOrderly()

Parâmetro

Descrição

Intervalo

numOfMessages

Quantidade máxima de mensagens a serem consumidas por lote

Até 16

waitSeconds

Duração do long polling em segundos. Na ausência de mensagens, a requisição fica retida no broker até a chegada de uma mensagem ou o fim do tempo limite

Até 30

Tratamento de erros

O consumidor lida com três cenários de erro:

Exceção

Causa

Ação

MessageResolveException

Caracteres inválidos no corpo da mensagem impedem a análise sintática

Confirmar (ACK) tanto as mensagens analisáveis quanto as não analisáveis

MessageNotExistException

Nenhuma mensagem disponível no tópico

Continuar o long polling automaticamente

AckMessageException

O receipt handle expirou antes do envio do ACK

Registrar o erro em log e tentar novamente na próxima entrega