Skip to content

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
csharp
// 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
csharp
// 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
csharp
// 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

SistemaTipoPersistênciaEscalabilidadeLatênciaComplexidadeMelhor Para
RabbitMQMessage BrokerSimMédiaBaixaMédiaMicroservices, routing complexo
KafkaEvent StreamingSimAltaMédiaAltaEvent sourcing, analytics
Azure Service BusCloud ServiceSimAltaBaixaBaixaAplicações Azure
Apache PulsarEvent StreamingSimMuito AltaBaixaAltaSistemas de alta performance
AWS SQS/SNSCloud ServiceSimAltaMédiaBaixaAplicaçõ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: true para 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

SistemaLocal de ArmazenamentoPersistênciaDurabilidadePerformanceEscalabilidade
RabbitMQMemória + DiscoConfigurávelMédiaMuito AltaMédia
KafkaDisco (Logs)SempreMuito AltaAltaMuito Alta
Azure Service BusAzure StorageSempreMuito AltaMédiaAlta
Apache PulsarBookKeeperSempreMuito AltaMuito AltaMuito Alta
AWS SQSAWS StorageSempreMuito AltaMédiaAlta
Redis Pub/SubMemóriaNuncaBaixaExtremamente AltaMédia

🔧 Implementação em .NET

Configuração de Dependências

csharp
// 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

csharp
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

csharp
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"