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 consome mensagens ordenadas em ordem estrita FIFO (first-in-first-out). Este tópico fornece exemplos de código para enviar e receber mensagens ordenadas com o SDK de cliente TCP para .NET.

Modos de ordenação

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

  • Mensagens globalmente ordenadas: Todas as mensagens de um tópico são publicadas e consumidas em ordem FIFO estrita.

  • Mensagens ordenadas por partição: As mensagens de um tópico são distribuídas entre partições com base na chave de fragmentação. Dentro de cada partição, o consumo ocorre em ordem FIFO estrita.

A chave de fragmentação é um campo que direciona as mensagens ordenadas para uma partição específica, diferenciando-se das chaves de mensagens normais.

Importante

O broker do ApsaraMQ for RocketMQ 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, o broker as ordenará pelo horário de chegada, o que pode divergir da ordem de negócios pretendida. Para garantir a ordenação, envie todas as mensagens relacionadas a partir de um único produtor ou thread.

Para obter mais informações, consulte Mensagens ordenadas.

Pré-requisitos

Antes de começar, verifique se você:

  • Baixou o SDK para .NET. Para obter mais informações, consulte Notas de versão

  • Preparou o ambiente. Para obter mais informações, consulte Preparar o ambiente

  • Crie instâncias, tópicos e grupos de consumidores no console do ApsaraMQ for RocketMQ. Para obter mais informações, consulte Criar recursos

  • Obteve um par de AccessKey para sua conta Alibaba Cloud. Para obter mais informações, consulte Criar um par de AccessKey

Enviar mensagens ordenadas

O exemplo de código a seguir envia mensagens ordenadas usando o SDK de cliente TCP para .NET. Como todas as mensagens compartilham a mesma chave de fragmentação, elas são roteadas para a mesma partição e consumidas na ordem de envio.

Para acessar o repositório completo de códigos, consulte Repositório de códigos do ApsaraMQ for RocketMQ.

Substitua os seguintes espaços reservados pelos valores reais:

Espaço reservado

Descrição

Exemplo

<your-group-id>

ID do grupo de produtores criado no console

GID_example

<your-topic>

Tópico criado no console

T_example_topic_name

<your-tcp-endpoint>

Endpoint TCP obtido na página Instance Details do console

NameSrv_Addr

<your-log-path>

Caminho local para logs do SDK

C://log

<your-sharding-key>

Chave de fragmentação que define o roteamento da partição

App-Test

using System;
using ons;

public class OrderProducerExampleForEx
{
    public OrderProducerExampleForEx()
    {
    }

    static void Main(string[] args) {
        // Configure the producer properties.
        ONSFactoryProperty factoryInfo = new ONSFactoryProperty();
        // Make sure that the environment variables ALIBABA_CLOUD_ACCESS_KEY_ID
        // and ALIBABA_CLOUD_ACCESS_KEY_SECRET are configured.
        // AccessKey ID for authentication.
        factoryInfo.setFactoryProperty(ONSFactoryProperty::AccessKey, getenv("ALIBABA_CLOUD_ACCESS_KEY_ID"));
        // AccessKey secret for authentication.
        factoryInfo.setFactoryProperty(ONSFactoryProperty::SecretKey, getenv("ALIBABA_CLOUD_ACCESS_KEY_SECRET"));
        // Producer group ID created in the ApsaraMQ for RocketMQ console.
        factoryInfo.setFactoryProperty(ONSFactoryProperty.ProducerId, "<your-group-id>");
        // Topic created in the ApsaraMQ for RocketMQ console.
        factoryInfo.setFactoryProperty(ONSFactoryProperty.PublishTopics, "<your-topic>");
        // TCP endpoint from the Instance Details page in the console.
        factoryInfo.setFactoryProperty(ONSFactoryProperty.NAMESRV_ADDR, "<your-tcp-endpoint>");
        // Local path for SDK logs.
        factoryInfo.setFactoryProperty(ONSFactoryProperty.LogPath, "<your-log-path>");

        // Create a producer instance.
        // Producer instances are thread-safe and can send messages to different topics.
        // In most cases, one producer instance per thread is sufficient.
        OrderProducer producer = ONSFactory.getInstance().createOrderProducer(factoryInfo);

        // Start the producer.
        producer.start();

        // Create a message with topic, tag, and body.
        Message msg = new Message(factoryInfo.getPublishTopics(), "tagA", "Example message body");
        // Define the sharding key. Messages with the same sharding key
        // are sent to the same partition and consumed in order.
        string shardingKey = "<your-sharding-key>";
        for (int i = 0; i < 32; i++) {
            try
            {
                SendResultONS sendResult = producer.send(msg, shardingKey);
                Console.WriteLine("send success {0}", sendResult.getMessageId());
            }
            catch (Exception ex)
            {
                Console.WriteLine("send failure{0}", ex.ToString());
            }
        }

        // Shut down the producer before exiting the thread.
        producer.shutdown();

    }
}

Receber mensagens ordenadas

O exemplo de código a seguir recebe mensagens ordenadas usando o SDK de cliente TCP para .NET. O consumidor processa as mensagens em ordem FIFO dentro de cada partição.

Substitua os espaços reservados abaixo pelos valores reais:

Espaço reservado

Descrição

Exemplo

<your-group-id>

ID do grupo de consumidores criado no console

GID_example

<your-topic>

Tópico criado no console

T_example_topic_name

<your-tcp-endpoint>

Endpoint TCP disponível na página Instance Details do console

NameSrv_Addr

<your-log-path>

Caminho local para logs do SDK

C://log

using System;
using System.Text;
using System.Threading;
using ons;

namespace demo
{

    public class MyMsgOrderListener : MessageOrderListener
    {
        public MyMsgOrderListener()
        {

        }

        ~MyMsgOrderListener()
        {
        }

        public override ons.OrderAction consume(Message value, ConsumeOrderContext context)
        {
            Byte[] text = Encoding.Default.GetBytes(value.getBody());
            Console.WriteLine(Encoding.UTF8.GetString(text));
            return ons.OrderAction.Success;
        }
    }

    class OrderConsumerExampleForEx
    {
        static void Main(string[] args)
        {
            // Configure the consumer properties.
            ONSFactoryProperty factoryInfo = new ONSFactoryProperty();
            // AccessKey ID for authentication.
            factoryInfo.setFactoryProperty(ONSFactoryProperty.AccessKey, "Your access key");
            // AccessKey secret for authentication.
            factoryInfo.setFactoryProperty(ONSFactoryProperty.SecretKey, "Your access secret");
            // Consumer group ID created in the ApsaraMQ for RocketMQ console.
            factoryInfo.setFactoryProperty(ONSFactoryProperty.ConsumerId, "<your-group-id>");
            // Topic created in the ApsaraMQ for RocketMQ console.
            factoryInfo.setFactoryProperty(ONSFactoryProperty.PublishTopics, "<your-topic>");
            // TCP endpoint from the Instance Details page in the console.
            factoryInfo.setFactoryProperty(ONSFactoryProperty.NAMESRV_ADDR, "<your-tcp-endpoint>");
            // Local path for SDK logs.
            factoryInfo.setFactoryProperty(ONSFactoryProperty.LogPath, "<your-log-path>");

            // Create a consumer instance.
            OrderConsumer consumer = ONSFactory.getInstance().createOrderConsumer(factoryInfo);

            // Subscribe to the topic with a wildcard tag filter.
            consumer.subscribe(factoryInfo.getPublishTopics(), "*", new MyMsgOrderListener());

            // Start the consumer.
            consumer.start();

            // Keep the main thread alive to continue receiving messages.
            Thread.Sleep(30000);

            // Shut down the consumer when it is no longer needed.
            consumer.shutdown();
        }
    }
}

Próximos passos