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