一、事件溯源概述
事件溯源(Event Sourcing)是一种架构模式,通过存储所有状态变更事件来记录系统状态。系统状态可以通过重放所有事件来重建,而不是存储当前状态。
二、事件溯源原理
2.1 传统模式vs事件溯源
graph TD
A[传统模式] --> B[存储当前状态]
B --> C[更新时覆盖]
C --> D[丢失历史]
E[事件溯源] --> F[存储所有事件]
F --> G[追加新事件]
G --> H[保留完整历史]
H --> I[重放事件重建状态]
2.2 事件溯源流程
sequenceDiagram
participant Client as 客户端
participant Aggregate as 聚合根
participant EventStore as 事件存储
participant ReadModel as 读模型
Client->>Aggregate: 命令
Aggregate->>EventStore: 加载历史事件
EventStore-->>Aggregate: 返回事件列表
Aggregate->>Aggregate: 应用事件重建状态
Aggregate->>Aggregate: 执行业务逻辑
Aggregate->>Aggregate: 生成新事件
Aggregate->>EventStore: 保存新事件
EventStore-->>Aggregate: 保存成功
EventStore->>ReadModel: 发布事件
ReadModel->>ReadModel: 更新读模型
2.3 事件溯源核心概念
| 概念 | 描述 | 特点 |
|---|---|---|
| 事件 | 状态变更的记录 | 不可变、可追溯 |
| 聚合根 | 业务逻辑的边界 | 事务一致性边界 |
| 事件存储 | 存储所有事件 | 追加写入、不可修改 |
| 快照 | 聚合状态的快照 | 加速状态重建 |
| 投影 | 从事件生成读模型 | 可重复、可并行 |
三、事件定义
3.1 事件接口
public interface IDomainEvent
{
Guid Id { get; }
Guid AggregateId { get; }
int Version { get; }
DateTime Timestamp { get; }
string EventType { get; }
}
public abstract class DomainEvent : IDomainEvent
{
public Guid Id { get; } = Guid.NewGuid();
public Guid AggregateId { get; protected set; }
public int Version { get; protected set; }
public DateTime Timestamp { get; } = DateTime.UtcNow;
public string EventType { get; protected set; }
protected DomainEvent()
{
EventType = GetType().Name;
}
}
public class OrderCreatedEvent : DomainEvent
{
public Guid UserId { get; }
public List<OrderItemData> Items { get; }
public string ShippingAddress { get; }
public decimal TotalAmount { get; }
public OrderCreatedEvent(Guid aggregateId, Guid userId,
List<OrderItemData> items, string shippingAddress, decimal totalAmount)
{
AggregateId = aggregateId;
UserId = userId;
Items = items;
ShippingAddress = shippingAddress;
TotalAmount = totalAmount;
}
}
public class OrderStatusChangedEvent : DomainEvent
{
public OrderStatus PreviousStatus { get; }
public OrderStatus NewStatus { get; }
public OrderStatusChangedEvent(Guid aggregateId, int version,
OrderStatus previousStatus, OrderStatus newStatus)
{
AggregateId = aggregateId;
Version = version;
PreviousStatus = previousStatus;
NewStatus = newStatus;
}
}
3.2 事件序列化
public class EventSerializer
{
public byte[] Serialize(IDomainEvent @event)
{
var envelope = new EventEnvelope
{
Id = @event.Id,
AggregateId = @event.AggregateId,
Version = @event.Version,
Timestamp = @event.Timestamp,
EventType = @event.EventType,
Data = JsonSerializer.Serialize(@event, @event.GetType())
};
return Encoding.UTF8.GetBytes(JsonSerializer.Serialize(envelope));
}
public IDomainEvent Deserialize(byte[] data)
{
var envelope = JsonSerializer.Deserialize<EventEnvelope>(data);
var eventType = Type.GetType(envelope.EventType);
return (IDomainEvent)JsonSerializer.Deserialize(envelope.Data, eventType);
}
}
public class EventEnvelope
{
public Guid Id { get; set; }
public Guid AggregateId { get; set; }
public int Version { get; set; }
public DateTime Timestamp { get; set; }
public string EventType { get; set; }
public string Data { get; set; }
}
四、聚合根设计
4.1 聚合根实现
public abstract class AggregateRoot
{
private List<IDomainEvent> _uncommittedEvents = new List<IDomainEvent>();
public Guid Id { get; protected set; }
public int Version { get; protected set; }
protected void Apply(IDomainEvent @event)
{
((dynamic)this).Apply(@event);
Version++;
}
protected void AddEvent(IDomainEvent @event)
{
_uncommittedEvents.Add(@event);
}
public List<IDomainEvent> GetUncommittedEvents()
{
var events = _uncommittedEvents.ToList();
_uncommittedEvents.Clear();
return events;
}
public void LoadFromHistory(List<IDomainEvent> events)
{
foreach (var @event in events)
{
Apply(@event);
}
}
}
public class Order : AggregateRoot
{
private Guid _userId;
private List<OrderItem> _items = new List<OrderItem>();
private OrderStatus _status;
private string _shippingAddress;
private decimal _totalAmount;
private Order() { }
public static Order Create(Guid userId, List<OrderItemDto> items,
List<Product> products, string shippingAddress)
{
var order = new Order { Id = Guid.NewGuid() };
var orderItems = items.Select(item =>
{
var product = products.First(p => p.Id == item.ProductId);
return new OrderItem(item.ProductId, product.Name, item.Quantity, product.Price);
}).ToList();
var totalAmount = orderItems.Sum(i => i.TotalPrice);
order.AddEvent(new OrderCreatedEvent(order.Id, userId, orderItems, shippingAddress, totalAmount));
return order;
}
public void UpdateStatus(OrderStatus newStatus)
{
if (!CanTransitionTo(newStatus))
throw new InvalidOperationException($"Cannot transition from {_status} to {newStatus}");
AddEvent(new OrderStatusChangedEvent(Id, Version + 1, _status, newStatus));
}
private void Apply(OrderCreatedEvent @event)
{
_userId = @event.UserId;
_items = @event.Items.Select(i => new OrderItem(i.ProductId, i.ProductName, i.Quantity, i.Price)).ToList();
_status = OrderStatus.Pending;
_shippingAddress = @event.ShippingAddress;
_totalAmount = @event.TotalAmount;
}
private void Apply(OrderStatusChangedEvent @event)
{
_status = @event.NewStatus;
}
private bool CanTransitionTo(OrderStatus newStatus)
{
// 状态转换逻辑
return true;
}
}
五、事件存储
5.1 事件存储设计
-- 事件存储表
CREATE TABLE events (
id UUID PRIMARY KEY,
aggregate_id UUID NOT NULL,
version INT NOT NULL,
event_type VARCHAR(100) NOT NULL,
data JSONB NOT NULL,
timestamp TIMESTAMP NOT NULL,
created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP,
UNIQUE INDEX idx_aggregate_version (aggregate_id, version),
INDEX idx_aggregate_id (aggregate_id),
INDEX idx_timestamp (timestamp)
);
-- 快照表
CREATE TABLE snapshots (
id UUID PRIMARY KEY,
aggregate_id UUID NOT NULL,
version INT NOT NULL,
data JSONB NOT NULL,
created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP,
UNIQUE INDEX idx_snapshot_aggregate (aggregate_id),
INDEX idx_snapshot_version (aggregate_id, version)
);
5.2 事件存储实现
public class EventStore : IEventStore
{
private readonly IDbConnection _connection;
private readonly EventSerializer _serializer;
public async Task AppendAsync(Guid aggregateId, int expectedVersion, List<IDomainEvent> events)
{
using (var transaction = await _connection.BeginTransactionAsync())
{
var currentVersion = await GetCurrentVersionAsync(aggregateId);
if (currentVersion != expectedVersion)
throw new ConcurrencyException(expectedVersion, currentVersion);
foreach (var @event in events)
{
@event.Version = currentVersion + 1;
currentVersion++;
await _connection.ExecuteAsync(
"INSERT INTO events (id, aggregate_id, version, event_type, data, timestamp) " +
"VALUES (@Id, @AggregateId, @Version, @EventType, @Data, @Timestamp)",
new
{
@event.Id,
@event.AggregateId,
@event.Version,
@event.EventType,
Data = _serializer.Serialize(@event),
@event.Timestamp
},
transaction
);
}
await transaction.CommitAsync();
}
}
public async Task<List<IDomainEvent>> GetEventsAsync(Guid aggregateId)
{
var rows = await _connection.QueryAsync<EventRow>(
"SELECT data FROM events WHERE aggregate_id = @AggregateId ORDER BY version",
new { AggregateId = aggregateId }
);
return rows.Select(row => _serializer.Deserialize(row.Data)).ToList();
}
public async Task<List<IDomainEvent>> GetEventsAsync(Guid aggregateId, int fromVersion)
{
var rows = await _connection.QueryAsync<EventRow>(
"SELECT data FROM events WHERE aggregate_id = @AggregateId AND version >= @FromVersion ORDER BY version",
new { AggregateId = aggregateId, FromVersion = fromVersion }
);
return rows.Select(row => _serializer.Deserialize(row.Data)).ToList();
}
private async Task<int> GetCurrentVersionAsync(Guid aggregateId)
{
var version = await _connection.QuerySingleOrDefaultAsync<int?>(
"SELECT MAX(version) FROM events WHERE aggregate_id = @AggregateId",
new { AggregateId = aggregateId }
);
return version ?? -1;
}
}
六、快照机制
6.1 快照原理
快照是聚合状态的保存点,用于加速状态重建:
graph TD
A[重建聚合状态] --> B{是否有快照}
B -->|是| C[加载快照]
B -->|否| D[从头开始]
C --> E[加载快照后的事件]
E --> F[应用事件]
F --> G[返回聚合]
D --> H[加载所有事件]
H --> F
6.2 快照实现
public class SnapshotService
{
private readonly IDbConnection _connection;
private readonly JsonSerializerOptions _serializerOptions;
public async Task SaveSnapshotAsync(AggregateRoot aggregate)
{
var snapshotData = JsonSerializer.Serialize(aggregate, _serializerOptions);
await _connection.ExecuteAsync(
"INSERT INTO snapshots (id, aggregate_id, version, data) " +
"VALUES (@Id, @AggregateId, @Version, @Data) " +
"ON CONFLICT (aggregate_id) DO UPDATE SET version = @Version, data = @Data",
new
{
Id = Guid.NewGuid(),
aggregate.Id,
aggregate.Version,
Data = snapshotData
}
);
}
public async Task<Snapshot> GetSnapshotAsync(Guid aggregateId)
{
var row = await _connection.QuerySingleOrDefaultAsync<SnapshotRow>(
"SELECT version, data FROM snapshots WHERE aggregate_id = @AggregateId",
new { AggregateId = aggregateId }
);
if (row == null)
return null;
return new Snapshot
{
Version = row.Version,
Data = row.Data
};
}
public async Task DeleteSnapshotAsync(Guid aggregateId)
{
await _connection.ExecuteAsync(
"DELETE FROM snapshots WHERE aggregate_id = @AggregateId",
new { AggregateId = aggregateId }
);
}
}
public class Snapshot
{
public int Version { get; set; }
public string Data { get; set; }
}
6.3 快照策略
public class SnapshotPolicy
{
private readonly SnapshotService _snapshotService;
private readonly int _snapshotInterval = 100;
public async Task CheckAndCreateSnapshotAsync(AggregateRoot aggregate)
{
if (aggregate.Version % _snapshotInterval == 0)
{
await _snapshotService.SaveSnapshotAsync(aggregate);
}
}
public async Task<AggregateRoot> LoadAggregateAsync<T>(Guid aggregateId) where T : AggregateRoot, new()
{
var snapshot = await _snapshotService.GetSnapshotAsync(aggregateId);
var aggregate = new T();
if (snapshot != null)
{
aggregate = JsonSerializer.Deserialize<T>(snapshot.Data);
var events = await _eventStore.GetEventsAsync(aggregateId, snapshot.Version + 1);
aggregate.LoadFromHistory(events);
}
else
{
var events = await _eventStore.GetEventsAsync(aggregateId);
aggregate.LoadFromHistory(events);
}
return aggregate;
}
}
七、投影
7.1 投影定义
投影是从事件生成读模型的过程:
public interface IProjection
{
Task HandleAsync(IDomainEvent @event);
Task ResetAsync();
}
public class OrderProjection : IProjection
{
private readonly IOrderReadRepository _readRepository;
public async Task HandleAsync(IDomainEvent @event)
{
switch (@event)
{
case OrderCreatedEvent createdEvent:
await HandleOrderCreated(createdEvent);
break;
case OrderStatusChangedEvent statusChangedEvent:
await HandleOrderStatusChanged(statusChangedEvent);
break;
}
}
private async Task HandleOrderCreated(OrderCreatedEvent @event)
{
var readModel = new OrderReadModel
{
Id = @event.AggregateId,
UserId = @event.UserId,
Status = (int)OrderStatus.Pending,
TotalAmount = @event.TotalAmount,
ItemCount = @event.Items.Count,
CreatedAt = @event.Timestamp
};
await _readRepository.InsertAsync(readModel);
}
private async Task HandleOrderStatusChanged(OrderStatusChangedEvent @event)
{
var readModel = await _readRepository.GetByIdAsync(@event.AggregateId);
if (readModel != null)
{
readModel.Status = (int)@event.NewStatus;
readModel.UpdatedAt = @event.Timestamp;
await _readRepository.UpdateAsync(readModel);
}
}
public async Task ResetAsync()
{
await _readRepository.ClearAsync();
var events = await _eventStore.GetAllEventsAsync();
foreach (var @event in events)
{
await HandleAsync(@event);
}
}
}
7.2 投影类型
| 类型 | 描述 | 特点 |
|---|---|---|
| 实时投影 | 事件发生时立即更新 | 低延迟 |
| 批处理投影 | 定时批量更新 | 高吞吐量 |
| 按需投影 | 查询时生成 | 数据新鲜 |
| 并行投影 | 多个投影并行处理 | 高性能 |
八、事件溯源与CQRS集成
8.1 CQRS+事件溯源架构
graph TD
A[客户端] --> B[命令接口]
A --> C[查询接口]
B --> D[命令处理程序]
D --> E[聚合根]
E --> F[事件存储]
F --> G[事件发布]
G --> H[投影]
H --> I[读模型]
C --> J[查询处理程序]
J --> I
8.2 命令处理流程
public class EventSourcingCommandHandler<TCommand, TAggregate>
: ICommandHandler<TCommand>
where TCommand : ICommand
where TAggregate : AggregateRoot, new()
{
private readonly IEventStore _eventStore;
private readonly SnapshotPolicy _snapshotPolicy;
private readonly IEventBus _eventBus;
public async Task<CommandResult> HandleAsync(TCommand command)
{
var aggregate = await LoadAggregateAsync(command.AggregateId);
await ApplyCommand(aggregate, command);
var events = aggregate.GetUncommittedEvents();
await _eventStore.AppendAsync(aggregate.Id, aggregate.Version - events.Count, events);
await _snapshotPolicy.CheckAndCreateSnapshotAsync(aggregate);
foreach (var @event in events)
{
await _eventBus.PublishAsync(@event);
}
return CommandResult.Success(aggregate.Id);
}
private async Task<TAggregate> LoadAggregateAsync(Guid aggregateId)
{
return await _snapshotPolicy.LoadAggregateAsync<TAggregate>(aggregateId);
}
protected abstract Task ApplyCommand(TAggregate aggregate, TCommand command);
}
九、事件溯源优势与挑战
9.1 优势
| 优势 | 描述 |
|---|---|
| 完整审计 | 所有变更都有记录 |
| 可追溯 | 可以追溯到任意时间点 |
| 可重放 | 可以重放事件重建状态 |
| 解耦 | 事件与投影解耦 |
| 灵活 | 可以添加新的投影 |
9.2 挑战
| 挑战 | 描述 | 解决方案 |
|---|---|---|
| 存储开销 | 事件数量随时间增长 | 快照、数据归档 |
| 重建时间 | 事件多重建慢 | 快照机制 |
| 并发冲突 | 乐观并发控制 | 版本号检查 |
| 事件版本 | 事件结构变更 | 事件版本化 |
十、事件溯源最佳实践
10.1 合理设计聚合
聚合应该是事务一致性的边界,不要设计过大的聚合。
10.2 使用快照
定期创建快照,避免重建时重放太多事件。
10.3 事件版本化
当事件结构变更时,使用版本化处理。
10.4 异步投影
使用消息队列异步处理投影,提高系统性能。
10.5 监控事件流
监控事件存储和投影的状态,及时发现问题。
十一、总结
事件溯源是一种强大的架构模式,通过存储所有状态变更事件,实现完整的审计追踪和状态重建能力。事件溯源与CQRS结合使用,能够构建高可扩展、高可维护的系统。实现事件溯源需要处理快照、投影、并发等复杂问题,适合有一定规模和复杂度的系统。