Sistemas de Mensageria para .NET Senior
📋 Visão Geral
Sistemas de mensageria são fundamentais para arquiteturas distribuídas, permitindo comunicação assíncrona entre serviços. Eles resolvem problemas de escalabilidade, resiliência e desacoplamento.
🏗️ Tipos de Sistemas de Mensageria
1. Message Brokers (Intermediários de Mensagem)
RabbitMQ
O que é: Message broker open-source que implementa AMQP (Advanced Message Queuing Protocol).
Características:
- Protocolo: AMQP, MQTT, STOMP
- Armazenamento: Em memória ou disco
- Padrão: Publish/Subscribe, Point-to-Point
- Persistência: Sim, com confirmações
- Escalabilidade: Horizontal com clustering
Quando usar:
- Comunicação entre microservices
- Workloads com diferentes padrões de mensagem
- Quando precisa de flexibilidade de routing
- Sistemas que precisam de garantias de entrega
// Producer
public class RabbitMQProducer
{
private readonly IConnection _connection;
private readonly IModel _channel;
public RabbitMQProducer(IConnectionFactory connectionFactory)
{
_connection = connectionFactory.CreateConnection();
_channel = _connection.CreateModel();
_channel.ExchangeDeclare("orders", ExchangeType.Topic);
_channel.QueueDeclare("order-created", durable: true, exclusive: false, autoDelete: false);
_channel.QueueBind("order-created", "orders", "order.created");
}
public void PublishOrderCreated(OrderCreatedEvent @event)
{
var message = JsonSerializer.Serialize(@event);
var body = Encoding.UTF8.GetBytes(message);
_channel.BasicPublish(
exchange: "orders",
routingKey: "order.created",
basicProperties: null,
body: body);
}
}
// Consumer
public class RabbitMQConsumer
{
private readonly IConnection _connection;
private readonly IModel _channel;
public RabbitMQConsumer(IConnectionFactory connectionFactory)
{
_connection = connectionFactory.CreateConnection();
_channel = _connection.CreateModel();
_channel.BasicQos(prefetchSize: 0, prefetchCount: 1, global: false);
}
public void StartConsuming()
{
var consumer = new EventingBasicConsumer(_channel);
consumer.Received += (model, ea) =>
{
var body = ea.Body.ToArray();
var message = Encoding.UTF8.GetString(body);
var @event = JsonSerializer.Deserialize<OrderCreatedEvent>(message);
ProcessOrderCreated(@event);
_channel.BasicAck(ea.DeliveryTag, false);
};
_channel.BasicConsume(queue: "order-created", autoAck: false, consumer: consumer);
}
}Apache Kafka
O que é: Plataforma de streaming distribuída para processamento de dados em tempo real.
Características:
- Protocolo: TCP customizado
- Armazenamento: Log distribuído persistente
- Padrão: Publish/Subscribe
- Persistência: Sim, com retenção configurável
- Escalabilidade: Alta, com particionamento
Quando usar:
- Processamento de streams de dados
- Event sourcing
- Analytics em tempo real
- Quando precisa de alta throughput
- Sistemas que precisam de replay de mensagens
// Producer
public class KafkaProducer
{
private readonly IProducer<string, string> _producer;
public KafkaProducer(ProducerConfig config)
{
_producer = new ProducerBuilder<string, string>(config).Build();
}
public async Task PublishOrderCreatedAsync(OrderCreatedEvent @event)
{
var message = JsonSerializer.Serialize(@event);
var result = await _producer.ProduceAsync(
topic: "order-events",
key: @event.OrderId.ToString(),
value: message);
Console.WriteLine($"Message delivered to {result.TopicPartitionOffset}");
}
}
// Consumer
public class KafkaConsumer
{
private readonly IConsumer<string, string> _consumer;
public KafkaConsumer(ConsumerConfig config)
{
_consumer = new ConsumerBuilder<string, string>(config).Build();
_consumer.Subscribe("order-events");
}
public void StartConsuming()
{
try
{
while (true)
{
var result = _consumer.Consume();
var @event = JsonSerializer.Deserialize<OrderCreatedEvent>(result.Message.Value);
ProcessOrderCreated(@event);
}
}
catch (OperationCanceledException)
{
_consumer.Close();
}
}
}Azure Service Bus
O que é: Serviço de mensageria gerenciado da Microsoft na nuvem.
Características:
- Protocolo: AMQP, HTTP/REST
- Armazenamento: Gerenciado pela Microsoft
- Padrão: Queues, Topics/Subscriptions
- Persistência: Sim, com TTL configurável
- Escalabilidade: Automática
Quando usar:
- Aplicações na nuvem Microsoft
- Quando não quer gerenciar infraestrutura
- Integração com outros serviços Azure
- Sistemas que precisam de alta disponibilidade
// Producer
public class ServiceBusProducer
{
private readonly ServiceBusClient _client;
private readonly ServiceBusSender _sender;
public ServiceBusProducer(string connectionString)
{
_client = new ServiceBusClient(connectionString);
_sender = _client.CreateSender("order-events");
}
public async Task PublishOrderCreatedAsync(OrderCreatedEvent @event)
{
var message = JsonSerializer.Serialize(@event);
var serviceBusMessage = new ServiceBusMessage(message)
{
MessageId = @event.OrderId.ToString(),
CorrelationId = @event.OrderId.ToString()
};
await _sender.SendMessageAsync(serviceBusMessage);
}
}
// Consumer
public class ServiceBusConsumer
{
private readonly ServiceBusClient _client;
private readonly ServiceBusProcessor _processor;
public ServiceBusConsumer(string connectionString)
{
_client = new ServiceBusClient(connectionString);
_processor = _client.CreateProcessor("order-events");
_processor.ProcessMessageAsync += ProcessMessageAsync;
_processor.ProcessErrorAsync += ProcessErrorAsync;
}
public async Task StartProcessingAsync()
{
await _processor.StartProcessingAsync();
}
private async Task ProcessMessageAsync(ProcessMessageEventArgs args)
{
var message = args.Message.Body.ToString();
var @event = JsonSerializer.Deserialize<OrderCreatedEvent>(message);
await ProcessOrderCreatedAsync(@event);
await args.CompleteMessageAsync(args.Message);
}
}2. Event Streaming Platforms
Apache Pulsar
O que é: Plataforma de mensageria distribuída que combina streaming e queuing.
Características:
- Protocolo: Pulsar Protocol
- Armazenamento: Apache BookKeeper
- Padrão: Publish/Subscribe, Queuing
- Persistência: Sim, com múltiplas camadas
- Escalabilidade: Muito alta
Quando usar:
- Sistemas que precisam de alta performance
- Quando precisa de múltiplos padrões de mensagem
- Sistemas com alta latência
- Quando precisa de geo-replicação
3. Cloud-Native Message Services
AWS SQS/SNS
O que é: Serviços de mensageria gerenciados da AWS.
Características:
- SQS: Queues simples e FIFO
- SNS: Publish/Subscribe
- Armazenamento: Gerenciado pela AWS
- Escalabilidade: Automática
Quando usar:
- Aplicações na AWS
- Quando não quer gerenciar infraestrutura
- Sistemas simples de mensageria
📊 Comparação de Sistemas
| Sistema | Tipo | Persistência | Escalabilidade | Latência | Complexidade | Melhor Para |
|---|---|---|---|---|---|---|
| RabbitMQ | Message Broker | Sim | Média | Baixa | Média | Microservices, routing complexo |
| Kafka | Event Streaming | Sim | Alta | Média | Alta | Event sourcing, analytics |
| Azure Service Bus | Cloud Service | Sim | Alta | Baixa | Baixa | Aplicações Azure |
| Apache Pulsar | Event Streaming | Sim | Muito Alta | Baixa | Alta | Sistemas de alta performance |
| AWS SQS/SNS | Cloud Service | Sim | Alta | Média | Baixa | Aplicações AWS |
🗄️ Armazenamento de Dados
RabbitMQ
Como funciona:
- Memória: Mensagens são mantidas em memória para máxima performance
- Disco: Persistência opcional usando arquivos de log
- Configuração:
durable: truepara persistência em disco - Limpeza: TTL (Time-To-Live) configurável por mensagem ou fila
- Estrutura: Cada fila é um arquivo separado no disco
- Recuperação: Ao reiniciar, reconstrói estado das filas duráveis
Vantagens:
- Performance muito alta quando em memória
- Flexibilidade entre performance e durabilidade
- Configuração granular por fila
Desvantagens:
- Limitação de memória para grandes volumes
- Necessidade de balancear performance vs durabilidade
Kafka
Como funciona:
- Log Distribuído: Mensagens são anexadas ao final de logs imutáveis
- Particionamento: Cada tópico é dividido em partições (shards)
- Replicação: Cada partição tem múltiplas réplicas para durabilidade
- Retenção: Configurável por tópico (tempo ou tamanho)
- Compressão: Suporte a compressão para economizar espaço
- Segmentação: Logs são divididos em segmentos para limpeza
Vantagens:
- Escalabilidade horizontal ilimitada
- Durabilidade muito alta com replicação
- Performance linear com número de partições
- Retenção configurável por tópico
Desvantagens:
- Complexidade de configuração
- Overhead de replicação
- Necessidade de Zookeeper (até versão 2.8)
Azure Service Bus
Como funciona:
- Azure Storage: Mensagens armazenadas no Azure Storage Account
- TTL: Time-to-live configurável por mensagem
- Dead Letter Queue: Mensagens que falharam são movidas para DLQ
- Sessões: Para garantir ordenação de mensagens
- Partitioning: Suporte a partições para alta disponibilidade
- Geo-Disaster Recovery: Replicação entre regiões
Vantagens:
- Gerenciado pela Microsoft
- Alta disponibilidade
- Integração nativa com Azure
- Recursos avançados (sessões, DLQ)
Desvantagens:
- Vendor lock-in
- Custo pode ser alto para grandes volumes
- Latência maior que soluções on-premise
Apache Pulsar
Como funciona:
- Apache BookKeeper: Sistema de storage distribuído
- Separation of Concerns: Storage separado da computação
- Multi-tenancy: Suporte nativo a múltiplos tenants
- Geo-replication: Replicação entre datacenters
- Tiered Storage: Hot storage (SSD) + Cold storage (S3)
- Compression: Suporte a múltiplos algoritmos
Vantagens:
- Performance muito alta
- Escalabilidade horizontal
- Flexibilidade de deployment
- Suporte a múltiplos padrões
Desvantagens:
- Complexidade de operação
- Curva de aprendizado alta
- Menos documentação que Kafka
AWS SQS/SNS
Como funciona:
- SQS Standard: Mensagens processadas pelo menos uma vez
- SQS FIFO: Mensagens processadas exatamente uma vez, em ordem
- SNS: Publish/Subscribe sem armazenamento persistente
- Visibility Timeout: Tempo que mensagem fica invisível após leitura
- Message Deduplication: Para FIFO queues
- Dead Letter Queue: Configurável para mensagens que falharam
Vantagens:
- Totalmente gerenciado
- Escalabilidade automática
- Integração nativa com AWS
- Sem overhead de operação
Desvantagens:
- Vendor lock-in
- Limitações de throughput
- Custo pode ser alto
- Menos controle sobre configurações
Redis Pub/Sub
Como funciona:
- Memória: Mensagens mantidas apenas em memória
- Sem Persistência: Mensagens são perdidas se não processadas
- Channels: Sistema de canais para pub/sub
- Pattern Matching: Suporte a padrões de canais
- Clustering: Suporte a Redis Cluster
Vantagens:
- Performance extremamente alta
- Simplicidade de uso
- Baixa latência
- Suporte a padrões complexos
Desvantagens:
- Sem persistência de mensagens
- Limitação de memória
- Mensagens podem ser perdidas
📊 Comparação de Armazenamento
| Sistema | Local de Armazenamento | Persistência | Durabilidade | Performance | Escalabilidade |
|---|---|---|---|---|---|
| RabbitMQ | Memória + Disco | Configurável | Média | Muito Alta | Média |
| Kafka | Disco (Logs) | Sempre | Muito Alta | Alta | Muito Alta |
| Azure Service Bus | Azure Storage | Sempre | Muito Alta | Média | Alta |
| Apache Pulsar | BookKeeper | Sempre | Muito Alta | Muito Alta | Muito Alta |
| AWS SQS | AWS Storage | Sempre | Muito Alta | Média | Alta |
| Redis Pub/Sub | Memória | Nunca | Baixa | Extremamente Alta | Média |
🔧 Implementação em .NET
Configuração de Dependências
// RabbitMQ
services.AddSingleton<IConnectionFactory>(provider =>
{
return new ConnectionFactory
{
HostName = "localhost",
UserName = "guest",
Password = "guest"
};
});
// Kafka
services.AddSingleton<ProducerConfig>(provider =>
{
return new ProducerConfig
{
BootstrapServers = "localhost:9092"
};
});
// Azure Service Bus
services.AddSingleton<ServiceBusClient>(provider =>
{
return new ServiceBusClient(connectionString);
});Padrões de Implementação
Outbox Pattern
public class OrderService
{
private readonly IOrderRepository _orderRepository;
private readonly IOutboxRepository _outboxRepository;
public async Task<Order> CreateOrderAsync(CreateOrderRequest request)
{
using var transaction = await _orderRepository.BeginTransactionAsync();
try
{
var order = new Order(request.Items);
await _orderRepository.SaveAsync(order);
var @event = new OrderCreatedEvent(order.Id, order.TotalAmount);
await _outboxRepository.SaveAsync(@event);
await transaction.CommitAsync();
return order;
}
catch
{
await transaction.RollbackAsync();
throw;
}
}
}Event Sourcing
public class EventStore
{
private readonly IEventRepository _eventRepository;
public async Task SaveEventsAsync(string aggregateId, IEnumerable<IDomainEvent> events, long expectedVersion)
{
var eventList = events.ToList();
var currentVersion = expectedVersion;
foreach (var @event in eventList)
{
currentVersion++;
@event.Version = currentVersion;
await _eventRepository.SaveAsync(aggregateId, @event);
}
}
public async Task<IEnumerable<IDomainEvent>> GetEventsAsync(string aggregateId)
{
return await _eventRepository.GetEventsAsync(aggregateId);
}
}⚠️ Considerações Importantes
Resiliência
- Retry Policies: Implementar retry com backoff exponencial
- Circuit Breaker: Proteger contra falhas em cascata
- Dead Letter Queues: Para mensagens que falharam
- Health Checks: Monitorar saúde dos sistemas
Performance
- Batching: Agrupar mensagens para melhor throughput
- Compression: Comprimir mensagens grandes
- Connection Pooling: Reutilizar conexões
- Async/Await: Usar operações assíncronas
Monitoramento
- Métricas: Throughput, latência, erro rate
- Logging: Logs estruturados para debugging
- Tracing: Distributed tracing para rastrear mensagens
- Alerting: Alertas para problemas críticos
📚 Recursos Adicionais
Ferramentas
- RabbitMQ Management: Interface web para RabbitMQ
- Kafka UI: Interface para gerenciar Kafka
- Azure Service Bus Explorer: Ferramenta para Azure Service Bus
Libraries .NET
- RabbitMQ.Client: Cliente oficial RabbitMQ
- Confluent.Kafka: Cliente Kafka
- Azure.Messaging.ServiceBus: Cliente Azure Service Bus
Livros Recomendados
- "Kafka: The Definitive Guide"
- "RabbitMQ in Action"
- "Building Event-Driven Microservices"