事件溯源与CQRS实战:构建高并发读写分离系统

深入讲解事件溯源(Event Sourcing)与CQRS(命令查询职责分离)的核心概念与实现方法,涵盖事件存储、事件回放、读写模型同步、最终一致性等关键技术,提供Go语言完整实战代码。

引言

事件溯源(Event Sourcing)和CQRS是构建复杂业务系统的强大模式。事件溯源将状态变更记录为事件序列,CQRS将读写操作分离优化。两者结合能够构建高并发、可审计、易扩展的系统。

核心概念

事件溯源(Event Sourcing)

传统CRUD模式:
┌─────────────────────────────────────────┐
│ UPDATE orders SET status='shipped'      │
│ WHERE id=123;                           │
│                                         │
│ 问题:状态被覆盖,无法追溯历史           │
└─────────────────────────────────────────┘

事件溯源模式:
┌─────────────────────────────────────────┐
│ 事件1: OrderCreated(id=123, user=456)   │
│ 事件2: OrderPaid(id=123, amount=100)    │
│ 事件3: OrderShipped(id=123, tracking=XYZ)│
│                                         │
│ 优势:完整历史、可回放、可审计           │
└─────────────────────────────────────────┘

CQRS(Command Query Responsibility Segregation)

传统模式(读写同一模型):
┌─────────┐         ┌──────────┐
│  Client  │ ──────▶ │ Database │
└─────────┘         └──────────┘

CQRS模式(读写分离):
┌─────────┐  写入   ┌──────────┐  事件   ┌──────────┐
│  Client  │ ──────▶ │ 写模型   │ ──────▶ │ 事件存储 │
└─────────┘         └──────────┘         └──────────┘
     ▲                                         │
     │ 读取                                    │ 投影
     │                                         ▼
┌─────────┐         ┌──────────┐         ┌──────────┐
│  Client  │ ◀────── │ 读模型   │ ◀────── │ 投影服务 │
└─────────┘         └──────────┘         └──────────┘

事件存储实现

事件定义

// 领域事件基类
type DomainEvent struct {
    EventID     string    `json:"event_id"`
    AggregateID string    `json:"aggregate_id"`
    EventType   string    `json:"event_type"`
    Timestamp   time.Time `json:"timestamp"`
    Version     int       `json:"version"`
    Metadata    map[string]string `json:"metadata,omitempty"`
    Payload     interface{} `json:"payload"`
}

// 订单领域事件
type OrderCreatedEvent struct {
    OrderID     string    `json:"order_id"`
    UserID      string    `json:"user_id"`
    Items       []OrderItem `json:"items"`
    TotalAmount Money     `json:"total_amount"`
    CreatedAt   time.Time `json:"created_at"`
}

type OrderPaidEvent struct {
    OrderID     string    `json:"order_id"`
    PaymentID   string    `json:"payment_id"`
    Amount      Money     `json:"amount"`
    PaidAt      time.Time `json:"paid_at"`
}

type OrderShippedEvent struct {
    OrderID       string    `json:"order_id"`
    TrackingNumber string   `json:"tracking_number"`
    Carrier       string    `json:"carrier"`
    ShippedAt     time.Time `json:"shipped_at"`
}

type OrderCancelledEvent struct {
    OrderID   string `json:"order_id"`
    Reason    string `json:"reason"`
    CancelledBy string `json:"cancelled_by"`
    CancelledAt time.Time `json:"cancelled_at"`
}

事件存储接口

type EventStore interface {
    // 保存事件到聚合根
    Save(ctx context.Context, aggregateID string, events []DomainEvent) error
    
    // 加载聚合根的所有事件
    Load(ctx context.Context, aggregateID string) ([]DomainEvent, error)
    
    // 从特定版本开始加载
    LoadFromVersion(ctx context.Context, aggregateID string, fromVersion int) ([]DomainEvent, error)
    
    // 订阅事件流(用于投影)
    Subscribe(ctx context.Context, eventTypes []string, handler EventHandler) error
}

type EventHandler func(ctx context.Context, event DomainEvent) error

PostgreSQL事件存储实现

type PostgresEventStore struct {
    db *sql.DB
}

func NewPostgresEventStore(db *sql.DB) *PostgresEventStore {
    return &PostgresEventStore{db: db}
}

// 创建事件表
/*
CREATE TABLE events (
    event_id UUID PRIMARY KEY,
    aggregate_id VARCHAR(255) NOT NULL,
    event_type VARCHAR(255) NOT NULL,
    version INT NOT NULL,
    timestamp TIMESTAMP NOT NULL DEFAULT NOW(),
    payload JSONB NOT NULL,
    metadata JSONB,
    
    UNIQUE(aggregate_id, version)
);

CREATE INDEX idx_events_aggregate ON events(aggregate_id, version);
CREATE INDEX idx_events_type ON events(event_type, timestamp);
*/

func (s *PostgresEventStore) Save(ctx context.Context, aggregateID string, events []DomainEvent) error {
    tx, err := s.db.BeginTx(ctx, nil)
    if err != nil {
        return err
    }
    defer tx.Rollback()
    
    stmt, err := tx.PrepareContext(ctx, `
        INSERT INTO events (event_id, aggregate_id, event_type, version, timestamp, payload, metadata)
        VALUES ($1, $2, $3, $4, $5, $6, $7)
    `)
    if err != nil {
        return err
    }
    defer stmt.Close()
    
    for _, event := range events {
        payloadJSON, _ := json.Marshal(event.Payload)
        metadataJSON, _ := json.Marshal(event.Metadata)
        
        _, err := stmt.ExecContext(ctx,
            event.EventID,
            event.AggregateID,
            event.EventType,
            event.Version,
            event.Timestamp,
            payloadJSON,
            metadataJSON,
        )
        
        if err != nil {
            // 处理版本冲突(并发写入)
            if isUniqueViolation(err) {
                return ErrConcurrencyConflict
            }
            return err
        }
    }
    
    return tx.Commit()
}

func (s *PostgresEventStore) Load(ctx context.Context, aggregateID string) ([]DomainEvent, error) {
    return s.LoadFromVersion(ctx, aggregateID, 0)
}

func (s *PostgresEventStore) LoadFromVersion(ctx context.Context, aggregateID string, fromVersion int) ([]DomainEvent, error) {
    rows, err := s.db.QueryContext(ctx, `
        SELECT event_id, aggregate_id, event_type, version, timestamp, payload, metadata
        FROM events
        WHERE aggregate_id = $1 AND version > $2
        ORDER BY version ASC
    `, aggregateID, fromVersion)
    
    if err != nil {
        return nil, err
    }
    defer rows.Close()
    
    var events []DomainEvent
    for rows.Next() {
        var event DomainEvent
        var payloadJSON, metadataJSON []byte
        
        err := rows.Scan(
            &event.EventID,
            &event.AggregateID,
            &event.EventType,
            &event.Version,
            &event.Timestamp,
            &payloadJSON,
            &metadataJSON,
        )
        if err != nil {
            return nil, err
        }
        
        // 反序列化payload
        event.Payload = s.deserializePayload(event.EventType, payloadJSON)
        json.Unmarshal(metadataJSON, &event.Metadata)
        
        events = append(events, event)
    }
    
    return events, nil
}

并发冲突处理

事件溯源系统中,多个请求同时操作同一聚合根时会产生并发冲突。通过乐观并发控制(Optimistic Concurrency Control) 和版本号机制解决:

// 版本号校验的 Save 方法
func (s *PostgresEventStore) SaveWithVersionCheck(
    ctx context.Context, 
    aggregateID string, 
    expectedVersion int,
    events []DomainEvent,
) error {
    tx, err := s.db.BeginTx(ctx, &sql.TxOptions{Isolation: sql.LevelSerializable})
    if err != nil {
        return err
    }
    defer tx.Rollback()
    
    // 1. 检查当前版本
    var currentVersion int
    err = tx.QueryRowContext(ctx, 
        "SELECT COALESCE(MAX(version), 0) FROM events WHERE aggregate_id = $1",
        aggregateID,
    ).Scan(&currentVersion)
    
    if err != nil {
        return err
    }
    
    // 2. 乐观锁:版本不匹配说明有并发写入
    if currentVersion != expectedVersion {
        return fmt.Errorf("concurrency conflict: expected version %d, found %d", 
            expectedVersion, currentVersion)
    }
    
    // 3. 分配版本号并保存
    for i, event := range events {
        event.Version = currentVersion + i + 1
        // 保存事件...
    }
    
    return tx.Commit()
}

// 客户端级别重试
func (h *OrderCommandHandler) handleWithRetry(
    ctx context.Context, 
    cmd CreateOrderCommand,
    maxRetries int,
) error {
    for i := 0; i < maxRetries; i++ {
        err := h.handleCreateOrder(ctx, cmd)
        if err == nil {
            return nil
        }
        
        // 如果是并发冲突,等待随机时间后重试
        if strings.Contains(err.Error(), "concurrency conflict") {
            jitter := time.Duration(rand.Intn(100)) * time.Millisecond
            time.Sleep(time.Duration(i*50)*time.Millisecond + jitter)
            continue
        }
        return err
    }
    return fmt.Errorf("max retries exceeded")
}

并发冲突解决策略对比:

策略适用场景优点缺点
乐观锁读多写少无锁开销、吞吐量高冲突时需要重试
悲观锁写冲突频繁避免重试、即时拒绝降低并发性能
合并策略最终一致性高可用合并逻辑复杂
最后写入获胜允许覆盖简单可能丢失数据

事件版本化与演进

随着时间推移,事件结构不可避免地会演变。需要建立事件版本化以支持向后兼容:

// 事件版本化接口
type VersionedEvent interface {
    EventVersion() int  // 返回事件结构版本号
}

// 订单创建事件 v1(原始版本)
type OrderCreatedEventV1 struct {
    OrderID   string  `json:"order_id"`
    UserID    string  `json:"user_id"`
    Total     float64 `json:"total"`
    CreatedAt string  `json:"created_at"`
}

// 订单创建事件 v2(增加了优惠码字段)
type OrderCreatedEventV2 struct {
    OrderID      string  `json:"order_id"`
    UserID       string  `json:"user_id"`
    Total        float64 `json:"total"`
    CouponCode   string  `json:"coupon_code,omitempty"`
    DiscountRate float64 `json:"discount_rate,omitempty"`
    CreatedAt    string  `json:"created_at"`
}

// 升级器:将 v1 事件升级为 v2
func UpgradeOrderCreatedV1ToV2(v1 OrderCreatedEventV1) OrderCreatedEventV2 {
    return OrderCreatedEventV2{
        OrderID:   v1.OrderID,
        UserID:    v1.UserID,
        Total:     v1.Total,
        CreatedAt: v1.CreatedAt,
        // CouponCode 和 DiscountRate 使用默认值
    }
}

// 反序列化时根据版本号选择对应结构体
func (s *PostgresEventStore) deserializePayload(eventType string, data []byte) interface{} {
    switch eventType {
    case "OrderCreated":
        // 先检查 payload 中是否有 version 字段
        var meta struct {
            EventVersion int `json:"_event_version"`
        }
        json.Unmarshal(data, &meta)
        
        if meta.EventVersion >= 2 {
            var v2 OrderCreatedEventV2
            json.Unmarshal(data, &v2)
            return v2
        }
        
        var v1 OrderCreatedEventV1
        json.Unmarshal(data, &v1)
        return UpgradeOrderCreatedV1ToV2(v1)
    }
    return nil
}

聚合根实现

// 订单聚合根
type OrderAggregate struct {
    ID          string
    UserID      string
    Status      OrderStatus
    Items       []OrderItem
    TotalAmount Money
    Version     int
    
    uncommittedEvents []DomainEvent
}

func (o *OrderAggregate) Apply(event DomainEvent) {
    switch e := event.Payload.(type) {
    case OrderCreatedEvent:
        o.ID = e.OrderID
        o.UserID = e.UserID
        o.Items = e.Items
        o.TotalAmount = e.TotalAmount
        o.Status = OrderStatusCreated
        
    case OrderPaidEvent:
        o.Status = OrderStatusPaid
        
    case OrderShippedEvent:
        o.Status = OrderStatusShipped
        
    case OrderCancelledEvent:
        o.Status = OrderStatusCancelled
    }
    
    o.Version = event.Version
}

// 创建订单(命令)
func (o *OrderAggregate) CreateOrder(cmd CreateOrderCommand) error {
    if o.ID != "" {
        return errors.New("order already exists")
    }
    
    if len(cmd.Items) == 0 {
        return errors.New("order must have at least one item")
    }
    
    event := OrderCreatedEvent{
        OrderID:     cmd.OrderID,
        UserID:      cmd.UserID,
        Items:       cmd.Items,
        TotalAmount: calculateTotal(cmd.Items),
        CreatedAt:   time.Now(),
    }
    
    o.recordEvent("OrderCreated", event)
    return nil
}

// 支付订单
func (o *OrderAggregate) Pay(cmd PayOrderCommand) error {
    if o.Status != OrderStatusCreated {
        return errors.New("order cannot be paid in current status")
    }
    
    if cmd.Amount != o.TotalAmount {
        return errors.New("payment amount mismatch")
    }
    
    event := OrderPaidEvent{
        OrderID:   o.ID,
        PaymentID: cmd.PaymentID,
        Amount:    cmd.Amount,
        PaidAt:    time.Now(),
    }
    
    o.recordEvent("OrderPaid", event)
    return nil
}

// 发货
func (o *OrderAggregate) Ship(cmd ShipOrderCommand) error {
    if o.Status != OrderStatusPaid {
        return errors.New("order must be paid before shipping")
    }
    
    event := OrderShippedEvent{
        OrderID:        o.ID,
        TrackingNumber: cmd.TrackingNumber,
        Carrier:        cmd.Carrier,
        ShippedAt:      time.Now(),
    }
    
    o.recordEvent("OrderShipped", event)
    return nil
}

func (o *OrderAggregate) recordEvent(eventType string, payload interface{}) {
    o.Version++
    
    event := DomainEvent{
        EventID:     uuid.New().String(),
        AggregateID: o.ID,
        EventType:   eventType,
        Version:     o.Version,
        Timestamp:   time.Now(),
        Payload:     payload,
    }
    
    o.uncommittedEvents = append(o.uncommittedEvents, event)
    o.Apply(event)
}

func (o *OrderAggregate) GetUncommittedEvents() []DomainEvent {
    return o.uncommittedEvents
}

func (o *OrderAggregate) ClearUncommittedEvents() {
    o.uncommittedEvents = nil
}

命令处理(Command Handler)

type OrderCommandHandler struct {
    eventStore EventStore
    eventBus   EventBus
}

func (h *OrderCommandHandler) HandleCreateOrder(ctx context.Context, cmd CreateOrderCommand) error {
    // 1. 加载聚合根(从事件重建状态)
    aggregate := &OrderAggregate{}
    
    events, err := h.eventStore.Load(ctx, cmd.OrderID)
    if err != nil {
        return err
    }
    
    for _, event := range events {
        aggregate.Apply(event)
    }
    
    // 2. 执行业务逻辑
    if err := aggregate.CreateOrder(cmd); err != nil {
        return err
    }
    
    // 3. 保存新事件
    uncommittedEvents := aggregate.GetUncommittedEvents()
    if err := h.eventStore.Save(ctx, cmd.OrderID, uncommittedEvents); err != nil {
        return err
    }
    
    // 4. 发布事件(异步投影)
    for _, event := range uncommittedEvents {
        h.eventBus.Publish(ctx, event)
    }
    
    aggregate.ClearUncommittedEvents()
    return nil
}

查询模型与投影

读模型定义

// 订单查询模型(优化读取)
type OrderReadModel struct {
    OrderID        string    `json:"order_id"`
    UserID         string    `json:"user_id"`
    Status         string    `json:"status"`
    TotalAmount    float64   `json:"total_amount"`
    Currency       string    `json:"currency"`
    ItemCount      int       `json:"item_count"`
    TrackingNumber string    `json:"tracking_number,omitempty"`
    CreatedAt      time.Time `json:"created_at"`
    UpdatedAt      time.Time `json:"updated_at"`
}

// 用户订单统计(聚合视图)
type UserOrderStats struct {
    UserID        string  `json:"user_id"`
    TotalOrders   int     `json:"total_orders"`
    TotalSpent    float64 `json:"total_spent"`
    PendingOrders int     `json:"pending_orders"`
    LastOrderAt   time.Time `json:"last_order_at"`
}

投影服务

type OrderProjectionService struct {
    db *sql.DB
}

func (s *OrderProjectionService) Project(ctx context.Context, event DomainEvent) error {
    switch e := event.Payload.(type) {
    case OrderCreatedEvent:
        return s.projectOrderCreated(ctx, e)
    case OrderPaidEvent:
        return s.projectOrderPaid(ctx, e)
    case OrderShippedEvent:
        return s.projectOrderShipped(ctx, e)
    case OrderCancelledEvent:
        return s.projectOrderCancelled(ctx, e)
    }
    
    return nil
}

func (s *OrderProjectionService) projectOrderCreated(ctx context.Context, event OrderCreatedEvent) error {
    _, err := s.db.ExecContext(ctx, `
        INSERT INTO order_read_models (
            order_id, user_id, status, total_amount, currency, 
            item_count, created_at, updated_at
        ) VALUES ($1, $2, $3, $4, $5, $6, $7, $7)
    `,
        event.OrderID,
        event.UserID,
        "created",
        event.TotalAmount.Amount,
        event.TotalAmount.Currency,
        len(event.Items),
        event.CreatedAt,
    )
    
    return err
}

func (s *OrderProjectionService) projectOrderPaid(ctx context.Context, event OrderPaidEvent) error {
    _, err := s.db.ExecContext(ctx, `
        UPDATE order_read_models 
        SET status = 'paid', updated_at = $2
        WHERE order_id = $1
    `, event.OrderID, event.PaidAt)
    
    return err
}

func (s *OrderProjectionService) projectOrderShipped(ctx context.Context, event OrderShippedEvent) error {
    _, err := s.db.ExecContext(ctx, `
        UPDATE order_read_models 
        SET status = 'shipped', tracking_number = $2, updated_at = $3
        WHERE order_id = $1
    `, event.OrderID, event.TrackingNumber, event.ShippedAt)
    
    return err
}

// 用户统计投影
type UserStatsProjectionService struct {
    db *sql.DB
}

func (s *UserStatsProjectionService) Project(ctx context.Context, event DomainEvent) error {
    switch e := event.Payload.(type) {
    case OrderCreatedEvent:
        return s.projectOrderCreated(ctx, e)
    case OrderPaidEvent:
        return s.projectOrderPaid(ctx, e)
    }
    return nil
}

func (s *UserStatsProjectionService) projectOrderCreated(ctx context.Context, event OrderCreatedEvent) error {
    _, err := s.db.ExecContext(ctx, `
        INSERT INTO user_order_stats (user_id, total_orders, pending_orders, last_order_at)
        VALUES ($1, 1, 1, $2)
        ON CONFLICT (user_id) DO UPDATE SET
            total_orders = user_order_stats.total_orders + 1,
            pending_orders = user_order_stats.pending_orders + 1,
            last_order_at = $2
    `, event.UserID, event.CreatedAt)
    
    return err
}

事件回放

重建聚合根状态

type OrderRepository struct {
    eventStore EventStore
}

func (r *OrderRepository) GetByID(ctx context.Context, orderID string) (*OrderAggregate, error) {
    events, err := r.eventStore.Load(ctx, orderID)
    if err != nil {
        return nil, err
    }
    
    if len(events) == 0 {
        return nil, ErrOrderNotFound
    }
    
    // 从事件重建状态
    aggregate := &OrderAggregate{}
    for _, event := range events {
        aggregate.Apply(event)
    }
    
    return aggregate, nil
}

完整事件回放(数据迁移)

type EventReplayer struct {
    eventStore EventStore
    projections []ProjectionService
}

func (r *EventReplayer) ReplayAll(ctx context.Context) error {
    // 清空读模型
    for _, projection := range r.projections {
        projection.Reset(ctx)
    }
    
    // 加载所有事件(按时间顺序)
    events, err := r.loadAllEvents(ctx)
    if err != nil {
        return err
    }
    
    log.Infof("Replaying %d events", len(events))
    
    // 重放所有事件
    for i, event := range events {
        for _, projection := range r.projections {
            if err := projection.Project(ctx, event); err != nil {
                log.Errorf("Failed to project event %d: %v", i, err)
                return err
            }
        }
        
        if (i+1) % 10000 == 0 {
            log.Infof("Replayed %d/%d events", i+1, len(events))
        }
    }
    
    log.Info("Event replay completed")
    return nil
}

最终一致性处理

事件总线

type EventBus interface {
    Publish(ctx context.Context, event DomainEvent) error
    Subscribe(ctx context.Context, eventType string, handler EventHandler) error
}

// Kafka事件总线实现
type KafkaEventBus struct {
    producer *kafka.Producer
    consumer *kafka.Consumer
    topic    string
}

func (b *KafkaEventBus) Publish(ctx context.Context, event DomainEvent) error {
    eventJSON, _ := json.Marshal(event)
    
    msg := &kafka.Message{
        Topic: b.topic,
        Key:   []byte(event.AggregateID),
        Value: eventJSON,
        Headers: []kafka.Header{
            {Key: "event_type", Value: []byte(event.EventType)},
        },
    }
    
    return b.producer.Produce(ctx, msg, nil)
}

func (b *KafkaEventBus) Subscribe(ctx context.Context, eventType string, handler EventHandler) error {
    go func() {
        for {
            msg, err := b.consumer.ReadMessage(ctx)
            if err != nil {
                log.Errorf("Failed to read message: %v", err)
                continue
            }
            
            // 检查事件类型
            msgEventType := getHeaderValue(msg, "event_type")
            if msgEventType != eventType {
                continue
            }
            
            var event DomainEvent
            json.Unmarshal(msg.Value, &event)
            
            if err := handler(ctx, event); err != nil {
                log.Errorf("Failed to handle event: %v", err)
                // 发送到死信队列
            }
        }
    }()
    
    return nil
}

总结

事件溯源与CQRS的核心价值:

事件溯源:

  • 完整历史:所有状态变更都有记录
  • 可审计:满足合规要求
  • 时间旅行:可以回溯到任意时间点的状态
  • 灵活性:可以从事件派生任意读模型

CQRS:

  • 读写分离:分别优化读写性能
  • 独立扩展:读模型和写模型可以独立扩展
  • 简化查询:读模型针对查询优化

适用场景:

  • 需要完整审计日志
  • 读写负载差异大
  • 业务逻辑复杂
  • 需要时间旅行能力

不适用场景:

  • 简单CRUD应用
  • 团队缺乏经验
  • 对一致性要求极高(需要强一致)

延伸阅读

继续阅读

探索更多技术文章

浏览归档,发现更多关于系统设计、工具链和工程实践的内容。

全部文章 返回首页

「backend」更多文章

  1. 零信任安全架构:从边界防御到身份中心的安全范式
  2. 混沌工程实践:构建高可用系统的故障注入与弹性测试
  3. 流式数据处理:Kafka Streams与Flink实战指南