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.
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 |
|
|
ID do grupo de produtores criado no console |
GID_example |
|
|
Tópico criado no console |
T_example_topic_name |
|
|
Endpoint TCP obtido na página Instance Details do console |
NameSrv_Addr |
|
|
Caminho local para logs do SDK |
C://log |
|
|
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 |
|
|
ID do grupo de consumidores criado no console |
GID_example |
|
|
Tópico criado no console |
T_example_topic_name |
|
|
Endpoint TCP disponível na página Instance Details do console |
NameSrv_Addr |
|
|
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
Saiba mais sobre os conceitos de mensagens ordenadas: Mensagens ordenadas
Explore outros exemplos no Repositório de códigos do ApsaraMQ for RocketMQ