📖 数据密集型设计

事件溯源架构模式

深入探讨事件溯源原理与实现策略

一、事件溯源概述

事件溯源(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结合使用,能够构建高可扩展、高可维护的系统。实现事件溯源需要处理快照、投影、并发等复杂问题,适合有一定规模和复杂度的系统。