CQRS 与事件溯源

CQRS 读写分离架构、事件存储机制、快照优化与投影设计,及一致性保障策略。

CQRS 与事件溯源

CQRS(Command Query Responsibility Segregation)将读写操作分离,事件溯源以事件序列替代状态快照。两者常结合使用,解决高并发读写冲突。

1. CQRS 核心思想

传统模式的问题

                    ┌─────────────┐
         ┌──读────→ │   统一模型   │
         │          │ (Entity)    │
         │          └─────────────┘
         │                  ↑
Client ──┤                  │
         │                  │
         └──写────→  数据库事务阻塞读

读写操作在同一模型上竞争资源,复杂查询需要大量 JOIN。

CQRS 分离

         ┌─────────────┐              ┌─────────────┐
 ┌──读──→ │  Query Side  │              │  Read Model │
 │        │ (简单、快速)  │ ←────────── │ (为查询优化)│
 │        └─────────────┘   同步/异步    └─────────────┘
Client                          ↑
 │         ┌─────────────┐     │
 └──写───→ │ Command Side│ ────┘  事件
           │ (业务校验)   │     事件溯源
           └─────────────┘
                    ↓
              ┌───────────┐
              │ Event Store│
              └───────────┘
维度CommandQuery
职责业务规则、状态变更数据检索
模型领域模型视图模型/投影
存储写库读库
优化一致性、事务性能、灵活查询

2. 事件溯源(Event Sourcing)

不存储当前状态,而是存储导致状态变更的所有事件

状态重建

Initial: Account(balance=0)

Events:
  [0] AccountCreated(id=123, holder="Alice")
  [1] Deposited(amount=100)
  [2] Withdrawn(amount=30)
  [3] Deposited(amount=50)
  
Rebuild: fold(apply event)
  → balance = 100 - 30 + 50 = 120

数据结构

public interface DomainEvent {
    UUID getAggregateId();
    long getVersion();
    Instant getOccurredOn();
}

public class BankAccount extends AggregateRoot {
    private UUID id;
    private BigDecimal balance;
    
    public void deposit(BigDecimal amount) {
        apply(new Deposited(id, amount, ++version));
    }
    
    public void withdraw(BigDecimal amount) {
        if (balance.compareTo(amount) < 0) 
            throw new InsufficientBalanceException();
        apply(new Withdrawn(id, amount, ++version));
    }
    
    @Override
    protected void when(DomainEvent event) {
        switch (event) {
            case Deposited e -> balance = balance.add(e.amount());
            case Withdrawn e -> balance = balance.subtract(e.amount());
        }
    }
}

3. 事件存储

-- 事件表设计
CREATE TABLE events (
    aggregate_id UUID,
    version INT,
    event_type VARCHAR(100),
    event_data JSONB,
    metadata JSONB,
    occurred_on TIMESTAMP,
    PRIMARY KEY (aggregate_id, version)
);

-- 应用事件(乐观锁)
INSERT INTO events (aggregate_id, version, ...)
VALUES (?, ?, ...)
ON CONFLICT (aggregate_id, version) DO NOTHING;
-- 影响行数为0 → 并发冲突

存储选型

方案特点
关系型数据库简单、事务支持好
EventStoreDB专用事件数据库,内置投影
Kafka高吞吐、保留策略

4. 快照优化

事件过多时重建状态耗时,引入快照:

每 N 个事件创建一次快照:

Events:  [1][2][3] ... [998][999][1000]
Snapshot at v1000: Account(balance=5000, version=1000)

重建时:加载 Snapshot + 应用 v1000 之后的事件
public class SnapshotRepository {
    @Transactional
    public <T extends AggregateRoot> T load(UUID id) {
        var snapshot = loadLatestSnapshot(id);
        var events = loadEventsAfter(id, snapshot.getVersion());
        return replay(snapshot, events);
    }
}

5. 投影(Projection)与读模型

// 监听器将事件投影到读模型
@Component
public class OrderProjection {
    @EventListener
    public void on(OrderCreatedEvent event) {
        readRepository.save(new OrderView(
            event.getOrderId(),
            event.getCustomerId(),
            event.getTotal()
        ));
    }
    
    @EventListener
    public void on(OrderPaidEvent event) {
        readRepository.updateStatus(event.getOrderId(), PAID);
    }
}

投影同步模式

模式延迟复杂度一致性
同步投影强一致
事务消息最终一致
异步消费者最终一致
CDC (Debezium)最终一致

6. 一致性保障

写模型一致性

  • 聚合内:单个聚合的事务一致性
  • 聚合间:最终一致性(通过领域事件)

读模型一致性

读模型滞后是正常情况。如何处理?

1. UI 乐观更新:前端先展示修改结果,后台异步同步
2. 读后写重试:写入后未在读取模型找到,延迟重试
3. 版本标记:写入后返回版本号,读模型校验版本

7. 适用场景

适用不适用
审计要求严格简单 CRUD
回溯/回放需求一致性强于性能
复杂业务规则团队无 DDD 经验
高并发写入短期项目

总结

概念含义
CQRS读写分离,独立优化
事件溯源以事件序列替代状态快照
投影事件转读模型
快照状态重建优化

CQRS + 事件溯源拥有强大的追溯能力和灵活查询能力,但引入了架构复杂度,需谨慎评估团队能力与业务需求。

继续阅读

探索更多技术文章

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

全部文章 返回首页

「架构」更多文章

  1. 架构评审与技术债管理
  2. SLA/SLO/SLI 与容量规划
  3. 云原生架构模式