Todos os produtos
Search
Central de documentação

ApsaraMQ for RocketMQ:Ordered messages

Última atualização: Jun 27, 2026

As mensagens ordenadas permitem que os consumidores processem as mensagens exatamente na ordem de envio. O ApsaraMQ for RocketMQ usa grupos de mensagens para definir o escopo da ordenação: as mensagens dentro do mesmo grupo são sempre entregues sequencialmente, enquanto as mensagens de grupos diferentes são processadas independentemente.

Casos de uso

Use mensagens ordenadas quando os sistemas downstream precisarem processar eventos na sequência exata em que ocorreram no upstream:

  • Correspondência de negociações -- No mercado financeiro, quando múltiplas ofertas compartilham o mesmo preço, aplica-se a regra de prioridade por ordem de chegada. O sistema de processamento de ordens deve tratar as ofertas na sequência exata de registro.

    Trade matchmaking

  • Sincronização incremental de dados -- Ao replicar alterações de banco de dados (inserções, atualizações, exclusões) por meio de uma fila de mensagens para um sistema de busca, reproduza as operações na ordem original. A reprodução fora de ordem gera um estado inconsistente.

    Normal message -- out-of-order risk

    Ordered messages -- consistent replay

Funcionamento da ordenação

A ordenação de mensagens ponta a ponta envolve duas partes: ordem de produção (lado do envio) e ordem de consumo (lado do recebimento). Ambas devem ser atendidas para garantir o processamento estrito FIFO (first in, first out).

Ordem de produção

Para garantir a ordem de produção, atenda às três condições abaixo:

  1. Mesmo grupo de mensagens -- Defina o mesmo grupo de mensagens para todas as mensagens que precisam ser ordenadas entre si. Mensagens em grupos diferentes não possuem relação de ordenação.

  2. Produtor único -- Envie todas as mensagens relacionadas a partir de um único produtor. Mesmo com o mesmo grupo de mensagens, envios de produtores diferentes em sistemas distintos não têm ordem determinística.

  3. Envio serial -- Envie as mensagens sequencialmente a partir de uma única thread. Embora o cliente produtor suporte acesso multithread, as mensagens enviadas em paralelo por threads diferentes não mantêm ordem determinística.

Quando essas condições são atendidas, as mensagens com o mesmo grupo de mensagens são armazenadas na mesma fila, respeitando a ordem de envio.

Storage logic for ordered messages

Comportamento de armazenamento:

  • As mensagens do mesmo grupo são armazenadas sequencialmente na mesma fila.

  • Mensagens de grupos diferentes podem coexistir na mesma fila, mas sua ordem relativa não é garantida.

No diagrama acima, o Grupo de Mensagens 1 (G1-M1, G1-M2, G1-M3) e o Grupo de Mensagens 4 (G4-M1, G4-M2) compartilham a Fila 1. O ApsaraMQ for RocketMQ garante a ordem dentro de cada grupo, mas não entre grupos distintos.

Ordem de consumo

Dois mecanismos trabalham em conjunto para garantir a ordem de consumo:

  1. Ordem de entrega -- O SDK e o protocolo do servidor entregam as mensagens na ordem de armazenamento. Siga rigorosamente o padrão receber-processar-confirmar. O processamento assíncrono pode quebrar a garantia de ordenação.

    Importante

    Com o PushConsumer, as mensagens são entregues uma por vez, seguindo a ordem de armazenamento. Com o SimpleConsumer, várias mensagens podem chegar em uma única operação de pull -- sua aplicação deve processá-las sequencialmente e confirmar cada uma antes de chamar receive novamente para o mesmo grupo. Para mais detalhes, consulte Tipos de consumidor.

  2. Retentativas limitadas -- Quando uma mensagem ordenada falha após atingir o número máximo de retentativas, ela é ignorada para desbloquear as mensagens subsequentes. Escolha um número de retentativas que equilibre a confiabilidade com o risco de bloquear todo o grupo. As retentativas dentro de um grupo de mensagens:

    • Não interrompem a garantia de ordenação.

    • Não afetam mensagens em outros grupos de mensagens.

    • Bloqueiam as mensagens subsequentes no mesmo grupo até que a mensagem atual seja resolvida.

    Importante

    Enquanto uma mensagem ordenada com falha está sendo reenviada, as mensagens subsequentes do mesmo grupo ficam bloqueadas. Elas só são entregues após a mensagem atual ser processada com sucesso ou esgotar suas tentativas de retentativa.

Combinações de ordem de produção e consumo

O modo FIFO estrito exige tanto a ordem de produção quanto a de consumo. No entanto, nem todo consumidor precisa de entrega ordenada. Combine as configurações conforme os requisitos de throughput e ordenação:

Ordem de produção

Ordem de consumo

Resultado

Grupo de mensagens definido; envio serial

Ordenada

FIFO estrito dentro de cada grupo de mensagens

Grupo de mensagens definido; envio serial

Concorrente

Melhor esforço de ordem cronológica; maior throughput

Sem grupo de mensagens; envio não ordenado

Ordenada

Ordenação estrita no nível da fila (segue a ordem de armazenamento, não a de envio)

Sem grupo de mensagens; envio não ordenado

Concorrente

Melhor esforço de ordem cronológica

Ciclo de vida da mensagem

Uma mensagem ordenada passa por cinco estados:

Message lifecycle

  1. Inicializada -- O produtor constrói a mensagem e prepara o envio.

  2. Pronta -- A mensagem chega ao broker e torna-se visível para os consumidores.

  3. Em trânsito -- Um consumidor recupera a mensagem e a processa. Se o broker não receber confirmação dentro do tempo limite, ele tenta entregar a mensagem novamente. Para mais detalhes, consulte Retentativa de consumo.

  4. Confirmada -- O consumidor confirma o resultado. Por padrão, o ApsaraMQ for RocketMQ retém todas as mensagens. O broker marca a mensagem como consumida, mas não a exclui imediatamente.

  5. Excluída -- Após o término do período de retenção ou quando o armazenamento estiver baixo, o broker exclui as mensagens mais antigas de forma rotativa. Antes da exclusão, as mensagens podem ser reconsumidas. Para mais detalhes, consulte Armazenamento e limpeza de mensagens.

Importante
  • Uma mensagem reenviada é tratada como uma nova mensagem. O ciclo de vida da mensagem original termina.

  • Enquanto uma mensagem ordenada está em processo de retentativa, as mensagens subsequentes do mesmo grupo permanecem bloqueadas até que a mensagem atual seja processada com sucesso.

Limites

  • Mensagens ordenadas só podem ser enviadas para tópicos com MessageType definido como FIFO. O tipo da mensagem deve corresponder ao tipo do tópico.

Importante

Se um grupo de consumidores estiver configurado para entrega ordenada, todas as mensagens consumidas por esse grupo serão faturadas como mensagens ordenadas, independentemente do tipo real. Caso a ordenação estrita não seja necessária, configure o grupo para entrega concorrente a fim de reduzir custos.

Pré-requisitos

Antes de começar, certifique-se de ter:

Importante

Mensagens ordenadas só podem ser enviadas para tópicos com MessageType definido como FIFO. Um grupo de consumidores que não esteja configurado para o modo de entrega ordenada entregará as mensagens de forma concorrente. Verifique tanto o tipo do tópico quanto o modo de entrega do grupo de consumidores antes de escrever qualquer código.

Otimizar a concorrência de consumo

Com o SDK gRPC do RocketMQ 5.x, o PushConsumer pode distribuir mensagens da mesma MessageQueue para threads diferentes com base no grupo de mensagens. Quanto mais distintos forem os valores dos seus grupos de mensagens, maior será o ganho de throughput.

Versões de SDK suportadas:

SDK

Versão mínima

Java

5.0.8

C++

5.0.3

Outros SDKs

Não suportado

Enviar e consumir mensagens ordenadas

Mensagens ordenadas exigem um grupo de mensagens em cada chamada de envio. Projete os grupos de mensagens com a granularidade mais fina que o seu negócio permitir -- por exemplo, use um ID de pedido ou ID de usuário. Isso permite a ordenação por entidade, maximizando o paralelismo entre entidades distintas.

Todos os exemplos abaixo usam Java. Para amostras completas de SDK em outras linguagens, consulte SDK gRPC do RocketMQ 5.x.

Código de exemplo

Enviar mensagens ordenadas

import org.apache.rocketmq.client.apis.*;
import org.apache.rocketmq.client.apis.message.Message;
import org.apache.rocketmq.client.apis.producer.Producer;
import org.apache.rocketmq.client.apis.producer.SendReceipt;

public class ProducerExample {
    public static void main(String[] args) throws ClientException {
        // Instance endpoint. Get it from the Endpoints tab on the instance details page.
        // Use the VPC endpoint for access from ECS over the internal network.
        // Use the public endpoint for access from a local machine or on-premises data center.
        String endpoints = "<your-endpoint>";
        // Topic name. Create the topic in the console first.
        String topic = "<your-fifo-topic>";
        ClientServiceProvider provider = ClientServiceProvider.loadService();
        ClientConfigurationBuilder builder = ClientConfiguration.newBuilder().setEndpoints(endpoints);
        // For public network access, set the instance username and password.
        // Get them from the Intelligent Authentication tab on the Access Control page.
        // For internal network access from ECS, skip this -- the server authenticates via VPC.
        // For Serverless instances, set credentials for public access.
        // If authentication-free internal access is enabled, skip this for internal access.
        // builder.setCredentialProvider(new StaticSessionCredentialsProvider("<instance-username>", "<instance-password>"));
        ClientConfiguration configuration = builder.build();
        Producer producer = provider.newProducerBuilder()
                .setTopics(topic)
                .setClientConfiguration(configuration)
                .build();
        // Build an ordered message with a message group.
        Message message = provider.newMessageBuilder()
                .setTopic(topic)
                .setKeys("messageKey")         // Message key for lookup
                .setTag("messageTag")          // Tag for consumer-side filtering
                .setMessageGroup("fifoGroup001") // Message group -- keep values discrete to avoid hot spots
                .setBody("messageBody".getBytes())
                .build();
        try {
            SendReceipt sendReceipt = producer.send(message);
            System.out.println(sendReceipt.getMessageId());
        } catch (ClientException e) {
            e.printStackTrace();
        }
    }
}

Substitua os seguintes espaços reservados pelos seus valores reais:

Espaço reservado

Descrição

Exemplo

<your-endpoint>

Endpoint da instância obtido no console

rmq-cn-xxx.cn-hangzhou.rmq.aliyuncs.com:8080

<your-fifo-topic>

Nome do tópico FIFO

order-events

<instance-username>

Nome de usuário da instância (apenas acesso público)

LTAI5tXxx

<instance-password>

Senha da instância (apenas acesso público)

xXxXxXx

Consumir com PushConsumer

O PushConsumer entrega mensagens ordenadas uma por vez, seguindo a ordem de armazenamento. Configure o grupo de consumidores para o modo de entrega ordenada no console antes de iniciar o consumidor. Caso contrário, as mensagens serão entregues de forma concorrente.

import org.apache.rocketmq.client.apis.*;
import org.apache.rocketmq.client.apis.consumer.ConsumeResult;
import org.apache.rocketmq.client.apis.consumer.FilterExpression;
import org.apache.rocketmq.client.apis.consumer.FilterExpressionType;
import org.apache.rocketmq.client.apis.consumer.PushConsumer;
import org.apache.rocketmq.shaded.org.slf4j.Logger;
import org.apache.rocketmq.shaded.org.slf4j.LoggerFactory;

import java.io.IOException;
import java.util.Collections;

public class PushConsumerExample {
    private static final Logger LOGGER = LoggerFactory.getLogger(PushConsumerExample.class);

    private PushConsumerExample() {
    }

    public static void main(String[] args) throws ClientException, IOException, InterruptedException {
        // Instance endpoint. Get it from the Endpoints tab on the instance details page.
        String endpoints = "<your-endpoint>";
        // Topic and consumer group. Create both in the console first.
        String topic = "<your-fifo-topic>";
        String consumerGroup = "<your-consumer-group>";
        final ClientServiceProvider provider = ClientServiceProvider.loadService();
        ClientConfigurationBuilder builder = ClientConfiguration.newBuilder().setEndpoints(endpoints);
        // For public network access, set credentials.
        // builder.setCredentialProvider(new StaticSessionCredentialsProvider("<instance-username>", "<instance-password>"));
        ClientConfiguration clientConfiguration = builder.build();
        // Subscribe to all tags.
        String tag = "*";
        FilterExpression filterExpression = new FilterExpression(tag, FilterExpressionType.TAG);
        PushConsumer pushConsumer = provider.newPushConsumerBuilder()
                .setClientConfiguration(clientConfiguration)
                .setConsumerGroup(consumerGroup)
                .setSubscriptionExpressions(Collections.singletonMap(topic, filterExpression))
                .setMessageListener(messageView -> {
                    // Process the message and return the result.
                    System.out.println("Consume Message: " + messageView);
                    return ConsumeResult.SUCCESS;
                })
                .build();
        Thread.sleep(Long.MAX_VALUE);
        // Close the consumer when no longer needed.
        // pushConsumer.close();
    }
}

Consumir com SimpleConsumer

O SimpleConsumer busca mensagens em lotes. Processe as mensagens dentro de cada lote sequencialmente e confirme cada uma individualmente.

Importante

Para mensagens no mesmo grupo de mensagens, se uma mensagem anterior não tiver sido confirmada, chamar receive novamente não retornará as mensagens subsequentes desse grupo. Esse mecanismo de bloqueio evita o processamento fora de ordem. Mensagens de outros grupos de mensagens não são afetadas e ainda podem ser recebidas simultaneamente.

import org.apache.rocketmq.client.apis.ClientConfiguration;
import org.apache.rocketmq.client.apis.ClientConfigurationBuilder;
import org.apache.rocketmq.client.apis.ClientException;
import org.apache.rocketmq.client.apis.ClientServiceProvider;
import org.apache.rocketmq.client.apis.consumer.FilterExpression;
import org.apache.rocketmq.client.apis.consumer.FilterExpressionType;
import org.apache.rocketmq.client.apis.consumer.SimpleConsumer;
import org.apache.rocketmq.client.apis.message.MessageView;
import org.apache.rocketmq.shaded.org.slf4j.Logger;
import org.apache.rocketmq.shaded.org.slf4j.LoggerFactory;

import java.io.IOException;
import java.time.Duration;
import java.util.Collections;
import java.util.List;

public class SimpleConsumerExample {
    private static final Logger LOGGER = LoggerFactory.getLogger(SimpleConsumerExample.class);

    private SimpleConsumerExample() {
    }

    public static void main(String[] args) throws ClientException, IOException {
        // Instance endpoint. Get it from the Endpoints tab on the instance details page.
        String endpoints = "<your-endpoint>";
        // Topic and consumer group. Create both in the console first.
        String topic = "<your-fifo-topic>";
        String consumerGroup = "<your-consumer-group>";
        final ClientServiceProvider provider = ClientServiceProvider.loadService();
        ClientConfigurationBuilder builder = ClientConfiguration.newBuilder().setEndpoints(endpoints);
        // For public network access, set credentials.
        // builder.setCredentialProvider(new StaticSessionCredentialsProvider("<instance-username>", "<instance-password>"));
        ClientConfiguration clientConfiguration = builder.build();

        Duration awaitDuration = Duration.ofSeconds(10);
        // Subscribe to all tags.
        String tag = "*";
        FilterExpression filterExpression = new FilterExpression(tag, FilterExpressionType.TAG);
        SimpleConsumer consumer = provider.newSimpleConsumerBuilder()
                .setClientConfiguration(clientConfiguration)
                .setConsumerGroup(consumerGroup)
                .setAwaitDuration(awaitDuration)                    // Long polling timeout
                .setSubscriptionExpressions(Collections.singletonMap(topic, filterExpression))
                .build();
        int maxMessageNum = 16;
        Duration invisibleDuration = Duration.ofSeconds(10);       // Invisibility window for processing
        // Poll for messages in a loop. Use multiple threads for higher throughput.
        while (true) {
            final List<MessageView> messageViewList = consumer.receive(maxMessageNum, invisibleDuration);
            messageViewList.forEach(messageView -> {
                System.out.println(messageView);
                // Acknowledge each message after processing.
                try {
                    consumer.ack(messageView);
                } catch (ClientException e) {
                    e.printStackTrace();
                }
            });
        }
        // Close the consumer when no longer needed.
        // consumer.close();
    }
}

Solucionar problemas de retentativas de consumo

As retentativas de consumo ordenado para o PushConsumer ocorrem no lado do cliente -- o servidor não registra detalhes das retentativas. Se um rastro de mensagem mostrar um resultado de entrega failed, verifique os logs do cliente consumidor.

Para o caminho do log do cliente, consulte Configuração de log.

Pesquise estas palavras-chave nos logs do cliente:

Message listener raised an exception while consuming messages
Failed to consume fifo message finally, run out of attempt times

Melhores práticas

Processe mensagens serialmente, não em lotes

Consuma uma mensagem por vez. O consumo em lote pode quebrar a ordenação.

Exemplo: As mensagens são enviadas na ordem 1 -> 2 -> 3 -> 4. Durante o consumo em lote, as mensagens 2 e 3 são processadas juntas e falham. Na retentativa, ambas 2 e 3 são reentregues -- mas a mensagem 2 pode ser processada novamente depois que a mensagem 3 já teve sucesso em outra tentativa, resultando em consumo fora de ordem.

Distribua grupos de mensagens para evitar pontos de congestionamento

O ApsaraMQ for RocketMQ usa o valor do grupo de mensagens para determinar qual fila no servidor armazena cada mensagem. Todas as mensagens do mesmo grupo são roteadas para a mesma fila. Concentrar muitas mensagens em poucos grupos sobrecarrega essas filas, criando pontos de congestionamento de armazenamento e limitando a escalabilidade.

Use chaves de granularidade fina como grupos de mensagens -- por exemplo, IDs de pedidos ou IDs de usuários. Quanto mais distintos forem os valores dos seus grupos de mensagens, mais uniformemente as mensagens serão distribuídas entre as filas. Isso mantém as mensagens da mesma entidade em ordem, enquanto distribui a carga entre filas para entidades diferentes.