Skip to content

Event Sourcing para .NET Senior

📋 Conceitos Fundamentais

O que é Event Sourcing?

Event Sourcing é um padrão arquitetural onde o estado de uma aplicação é determinado por uma sequência de eventos que ocorreram ao longo do tempo, em vez de armazenar apenas o estado atual.

Vantagens

  • Audit Trail Completo: Histórico completo de todas as mudanças
  • Temporal Queries: Capacidade de consultar o estado em qualquer ponto no tempo
  • Debugging: Facilita a identificação de problemas
  • Scalability: Permite otimizações específicas para leitura e escrita

🏗️ Implementação Prática

1. Domain Events

csharp
public abstract class DomainEvent
{
    public Guid Id { get; set; }
    public DateTime OccurredOn { get; set; }
    public long Version { get; set; }
    public string AggregateId { get; set; }
}

public class OrderCreatedEvent : DomainEvent
{
    public string CustomerId { get; set; }
    public decimal TotalAmount { get; set; }
    public List<OrderItem> Items { get; set; }
}

public class OrderItemAddedEvent : DomainEvent
{
    public string ProductId { get; set; }
    public int Quantity { get; set; }
    public decimal UnitPrice { get; set; }
}

2. Aggregate Root

csharp
public class Order : AggregateRoot
{
    private readonly List<OrderItem> _items = new();
    private OrderStatus _status;
    private string _customerId;
    private decimal _totalAmount;

    public string CustomerId => _customerId;
    public IReadOnlyList<OrderItem> Items => _items.AsReadOnly();
    public OrderStatus Status => _status;
    public decimal TotalAmount => _totalAmount;

    public Order(string customerId, List<OrderItem> items)
    {
        Apply(new OrderCreatedEvent
        {
            CustomerId = customerId,
            Items = items,
            TotalAmount = items.Sum(i => i.Quantity * i.UnitPrice)
        });
    }

    public void AddItem(string productId, int quantity, decimal unitPrice)
    {
        if (_status != OrderStatus.Draft)
            throw new InvalidOperationException("Cannot modify confirmed order");

        Apply(new OrderItemAddedEvent
        {
            ProductId = productId,
            Quantity = quantity,
            UnitPrice = unitPrice
        });
    }

    protected override void Apply(DomainEvent @event)
    {
        switch (@event)
        {
            case OrderCreatedEvent e:
                _customerId = e.CustomerId;
                _items.AddRange(e.Items);
                _totalAmount = e.TotalAmount;
                _status = OrderStatus.Draft;
                break;
            case OrderItemAddedEvent e:
                _items.Add(new OrderItem(e.ProductId, e.Quantity, e.UnitPrice));
                _totalAmount += e.Quantity * e.UnitPrice;
                break;
        }
    }
}

3. Event Store

csharp
public interface IEventStore
{
    Task<IEnumerable<DomainEvent>> GetEventsAsync(string aggregateId);
    Task SaveEventsAsync(string aggregateId, IEnumerable<DomainEvent> events, long expectedVersion);
}

public class SqlEventStore : IEventStore
{
    private readonly IDbConnection _connection;

    public async Task<IEnumerable<DomainEvent>> GetEventsAsync(string aggregateId)
    {
        var query = @"
            SELECT EventType, EventData, Version, OccurredOn
            FROM Events 
            WHERE AggregateId = @AggregateId 
            ORDER BY Version";

        var events = await _connection.QueryAsync<EventData>(query, new { AggregateId = aggregateId });
        return events.Select(DeserializeEvent);
    }

    public async Task SaveEventsAsync(string aggregateId, IEnumerable<DomainEvent> events, long expectedVersion)
    {
        using var transaction = _connection.BeginTransaction();
        
        try
        {
            var currentVersion = await GetCurrentVersionAsync(aggregateId);
            if (currentVersion != expectedVersion)
                throw new ConcurrencyException();

            foreach (var @event in events)
            {
                await InsertEventAsync(aggregateId, @event, currentVersion + 1);
                currentVersion++;
            }

            transaction.Commit();
        }
        catch
        {
            transaction.Rollback();
            throw;
        }
    }
}

4. Event Handlers

csharp
public class OrderEventHandler : IEventHandler<OrderCreatedEvent>, IEventHandler<OrderItemAddedEvent>
{
    private readonly IEmailService _emailService;
    private readonly IInventoryService _inventoryService;

    public async Task HandleAsync(OrderCreatedEvent @event)
    {
        await _emailService.SendOrderConfirmationAsync(@event.CustomerId, @event.TotalAmount);
    }

    public async Task HandleAsync(OrderItemAddedEvent @event)
    {
        await _inventoryService.ReserveInventoryAsync(@event.ProductId, @event.Quantity);
    }
}

📊 Projections e Read Models

1. Projection

csharp
public class OrderProjection : IProjection
{
    private readonly IDbConnection _connection;

    public async Task HandleAsync(OrderCreatedEvent @event)
    {
        var query = @"
            INSERT INTO OrderReadModel (Id, CustomerId, TotalAmount, Status, CreatedOn)
            VALUES (@Id, @CustomerId, @TotalAmount, @Status, @OccurredOn)";

        await _connection.ExecuteAsync(query, new
        {
            Id = @event.AggregateId,
            @event.CustomerId,
            @event.TotalAmount,
            Status = OrderStatus.Draft,
            @event.OccurredOn
        });
    }

    public async Task HandleAsync(OrderItemAddedEvent @event)
    {
        var query = @"
            UPDATE OrderReadModel 
            SET TotalAmount = TotalAmount + @AdditionalAmount
            WHERE Id = @OrderId";

        await _connection.ExecuteAsync(query, new
        {
            OrderId = @event.AggregateId,
            AdditionalAmount = @event.Quantity * @event.UnitPrice
        });
    }
}

🔄 Snapshots

Implementação de Snapshots

csharp
public class OrderSnapshot
{
    public string Id { get; set; }
    public string CustomerId { get; set; }
    public List<OrderItem> Items { get; set; }
    public OrderStatus Status { get; set; }
    public decimal TotalAmount { get; set; }
    public long Version { get; set; }
}

public class SnapshotService
{
    private readonly IEventStore _eventStore;
    private readonly ISnapshotStore _snapshotStore;
    private const int SnapshotFrequency = 100;

    public async Task<Order> LoadAggregateAsync(string aggregateId)
    {
        var snapshot = await _snapshotStore.GetLatestSnapshotAsync<OrderSnapshot>(aggregateId);
        
        if (snapshot != null)
        {
            var events = await _eventStore.GetEventsAsync(aggregateId, snapshot.Version + 1);
            return RebuildAggregateFromSnapshot(snapshot, events);
        }

        var allEvents = await _eventStore.GetEventsAsync(aggregateId);
        return RebuildAggregateFromEvents(allEvents);
    }

    public async Task SaveSnapshotAsync(string aggregateId, AggregateRoot aggregate, long version)
    {
        if (version % SnapshotFrequency == 0)
        {
            var snapshot = CreateSnapshot(aggregate, version);
            await _snapshotStore.SaveSnapshotAsync(aggregateId, snapshot);
        }
    }
}

🚀 Performance e Otimizações

1. Event Streaming

csharp
public class EventStreamProcessor
{
    private readonly IEventStore _eventStore;
    private readonly IEnumerable<IProjection> _projections;
    private readonly IMessageBroker _messageBroker;

    public async Task ProcessEventsAsync()
    {
        var events = await _eventStore.GetUnprocessedEventsAsync();
        
        foreach (var @event in events)
        {
            await ProcessEventAsync(@event);
            await _messageBroker.PublishAsync(@event);
        }
    }

    private async Task ProcessEventAsync(DomainEvent @event)
    {
        foreach (var projection in _projections)
        {
            await projection.HandleAsync(@event);
        }
    }
}

2. Caching Strategies

csharp
public class CachedEventStore : IEventStore
{
    private readonly IEventStore _innerEventStore;
    private readonly IDistributedCache _cache;

    public async Task<IEnumerable<DomainEvent>> GetEventsAsync(string aggregateId)
    {
        var cacheKey = $"events:{aggregateId}";
        var cached = await _cache.GetAsync(cacheKey);
        
        if (cached != null)
            return DeserializeEvents(cached);

        var events = await _innerEventStore.GetEventsAsync(aggregateId);
        await _cache.SetAsync(cacheKey, SerializeEvents(events), TimeSpan.FromMinutes(30));
        
        return events;
    }
}

🧪 Testes

1. Unit Tests

csharp
[Test]
public async Task When_OrderCreated_Then_OrderCreatedEventRaised()
{
    var customerId = "customer-123";
    var items = new List<OrderItem> { new("product-1", 2, 10.0m) };

    var order = new Order(customerId, items);

    var events = order.GetUncommittedEvents();
    Assert.That(events, Has.Exactly(1).Items);
    Assert.That(events.First(), Is.TypeOf<OrderCreatedEvent>());
}

[Test]
public async Task When_LoadOrderFromEvents_Then_OrderStateCorrect()
{
    var events = new List<DomainEvent>
    {
        new OrderCreatedEvent { CustomerId = "customer-123", TotalAmount = 20.0m },
        new OrderItemAddedEvent { ProductId = "product-2", Quantity = 1, UnitPrice = 15.0m }
    };

    var order = Order.FromEvents(events);

    Assert.That(order.CustomerId, Is.EqualTo("customer-123"));
    Assert.That(order.TotalAmount, Is.EqualTo(35.0m));
}

📚 Ferramentas e Frameworks

1. EventStoreDB

csharp
public class EventStoreDbEventStore : IEventStore
{
    private readonly EventStoreClient _client;

    public async Task<IEnumerable<DomainEvent>> GetEventsAsync(string aggregateId)
    {
        var streamName = $"order-{aggregateId}";
        var events = _client.ReadStreamAsync(streamName, StreamPosition.Start);
        
        var domainEvents = new List<DomainEvent>();
        await foreach (var @event in events)
        {
            var domainEvent = DeserializeEvent(@event.Event.Data);
            domainEvents.Add(domainEvent);
        }
        
        return domainEvents;
    }
}

2. Apache Kafka

csharp
public class KafkaEventPublisher : IEventPublisher
{
    private readonly IProducer<string, string> _producer;

    public async Task PublishAsync(DomainEvent @event)
    {
        var message = new Message<string, string>
        {
            Key = @event.AggregateId,
            Value = JsonSerializer.Serialize(@event)
        };

        await _producer.ProduceAsync("domain-events", message);
    }
}

⚠️ Considerações e Trade-offs

Vantagens

  • Audit Trail Completo: Histórico de todas as mudanças
  • Temporal Queries: Consultar estado em qualquer momento
  • Debugging: Facilita identificação de problemas
  • Scalability: Separação de leitura e escrita

Desvantagens

  • Complexidade: Aumenta a complexidade do sistema
  • Storage: Maior uso de armazenamento
  • Learning Curve: Curva de aprendizado íngreme
  • Performance: Overhead em operações simples

Quando Usar

  • Sistemas que precisam de audit trail completo
  • Aplicações com requisitos de compliance
  • Sistemas que precisam de temporal queries
  • Quando a escalabilidade é crítica

Quando Não Usar

  • CRUD simples sem requisitos de auditoria
  • Sistemas com restrições de performance
  • Equipes sem experiência em Event Sourcing
  • Quando a simplicidade é mais importante que funcionalidades avançadas