引言

事件溯源(Event Sourcing)和命令查询职责分离(CQRS)是两种相互促进的架构模式。事件溯源将系统的所有状态变化记录为不可变事件流,通过回放这些事件重建当前状态;CQRS则将写操作和读操作分别建模、独立优化。两者结合,在金融交易、实时协作、审计追踪、微服务协作等场景中展现出强大的工程价值。本文将深入剖析这两种模型的底层实现、典型陷阱与工程实践。

1. 事件溯源(ES):状态即事件之投影

传统存储是对"当前状态"的快照保存——更新即覆盖。而事件溯源记录的是"发生了什么"——【事件】是事实,【状态】只是事件的投影

以银行账户为例:

传统模式:account表 {balance: 150}
事件模式:
   AccountOpened   {accountId: A001, initial: 100} @t1
   MoneyDeposited  {accountId: A001, amount: 80}    @t2
   MoneyWithdrew   {accountId: A001, amount: 30}    @t3
  → 当前balance = 100 + 80 - 30 = 150

应用关心的是状态(150元),但存储的是事件流。事件本身不可变——一旦写入就不能修改。状态则是通过折叠(event folding)聚合(aggregation)操作投影出来的。

1.1 事件的三个核心特征

不可变性(Immutability):事件是过去发生的事实,一旦确认就不可修改。修正错误只能通过补偿事件(Compensating Event)——如MoneyDepositedExcessive伴随MoneyCorrection。

有序性(Ordering):同一聚合内事件严格有序。版本号(Version):每个事件携带全局递增版本号(或时间戳),防止并发更新冲突。

1.2 事件存储(Event Store)

事件存储是事件溯源的物理基础。本质是一个Append-only日志,保证事件的持久性和顺序性。核心能力包括:

  • Append:向流追加新事件,返回新版本号
  • Read:按offset/version读取一个或多个事件
  • Subscribe:订阅事件流变更(Change Feed/Event Notification)

主流Event Store实现:EventStoreDBAxon ServerApache Kafka(配合Log Compaction)、DynamoDB StreamsMySQL/PostgreSQL的事件表。事件存储的选择取决于是否需求原子聚合授权、持久订阅、snapshot内建等特性。

2. CQRS:命令与查询的分离

CQRS(Command Query Responsibility Segregation)是Command Query Separation(CQS)的演进。CQS在对象层面区分Command(修改,不返回结果)和Query(返回结果,不修改)。CQRS将这个思想放大到数据模型层面:写模型和读模型分离为两套独立优化的数据结构

CQRS的两个核心组件:

  • 写模型(Write Model):接收Commands,验证业务规则,产生Events,提交到事件存储
  • 读模型(Read Model):订阅Events投影到查询优化结构(各种视图)

读模型可以做水平扩展到多个视图:Table-UserSummary(快速展示用户摘要)、Table-UserDetail(完整详情)、Table-FriendsView(社交视图)。每个视图订阅事件流并构建自己的投影(projection)。

2.1 CQRS的优势

  • 读写独立优化:写入关注事务一致性与并发控制;读取关注访问模式与查询性能
  • 视图独立扩展:不同查询模式可独立部署、独立扩展(order_view部署10个实例,user_view部署2个)
  • 存储引擎分化:写入用关系型数据库保证ACID,读取用搜索引擎(全文)或NoSQL(高并发)
  • 性能隔离:复杂查询不影响写入吞吐

2.2 CQRS的代价

  • 最终一致性:读模型投影异步——写后立即可见(CQRS强一致) vs 写后可能延迟可见(CQRS最终一致)
  • 系统复杂度:双模型维护、投影一致性、事件回放bug修复
  • 事件回放风险:Bug-event被修复后,投影可能无法正确更新(需要重放或补偿)

3. EventSourcing + CQRS的组合

事件溯源与CQRS是天然组合:ES提供"持久化的状态变化流" + CQRS提供"读写模型分离+投影派生不同视图"。

3.1 写端(Write Side)流程

命令(Command)携带意图+参数被发送到CommandHandler,Handler加载聚合历史事件→重建聚合当前状态→执行业务决策→产出新事件(Event)。

Axon Framework的EventSourcing代码示例:

@Aggregate
public class AccountAggregate {
    @AggregateIdentifier
    private String accountId;
    private BigDecimal balance;

    @CommandHandler
    public AccountAggregate(OpenAccountCommand cmd) {
        apply(new AccountOpenedEvent(cmd.getAccountId(), cmd.getInitial()));
    }

    @EventSourcingHandler
    public void on(AccountOpenedEvent event) {
        this.accountId = event.getAccountId();
        this.balance = event.getInitial();
    }
}

CommandHandler构造聚合时不调用无参构造函数,而是调用@EventSourcingHandler方法按序还原。这就是事件溯源的独特工作方式。

3.2 读端(Read_side)投影

事件发布到EventBus(或Message Broker)后,投影处理器(Event Handler)注册订阅并更新物化视图:

@EventHandler
public void on(DepositedEvent event, @MetaDataValue("userId") String userId) {
    // 更新读模型视图
    accountSummaryRepository.addBalance(event.getAccountId(), event.getAmount());
}

单一事件类型可对应多个投影处理器。例如DepositedEvent被AccountSummaryProjection、AccountAuditProjection、NotificationService三个订阅方消费。

4. 快照(Snapshot)与聚合加载优化

事件溯源的核心性能问题:聚合加载需要回放所有历史事件。若一个账户有1亿条记录,每次命令触发1亿次回放,系统无法使用。

快照机制:定期(每100、1000、10000个事件)在eventstore中保存聚合当前状态。加载时先加载最新快照,再仅从该版本之后的事件回放。播放路径:快照v100 + 99n个事件(100~199)→ 当前状态。

写端建议快照策略:

  • 事件数量阈值:每100个事件存储快照(简单、稳定、预测性强)
  • 时间窗口:每5分钟聚合快照(适合活动周期性聚合)
  • 版本差异:快照&当前版本差异大到阈值存储(如事件大小超过总容量阈值)

Axon的Snapshot配置:

@Aggregate(snapshotFilter = "snapshotFilter", snapshotTriggerDefinition = "snapshotTriggerDef")
// 或SnapshotterConfiguration配置每个聚合的SnapshotTrigger

5. 事件版本化与演进

事件溯源的最大挑战之一是存储的事件结构随业务演化。今天保存的v1格式事件,明天业务需求变了需要v2格式。事件是不可变的——我们不能修改已有事件。

5.1 事件版本化典型策略

双事件更新(Upcasting):将旧事件格式在加载时动态升迁为新格式。每个旧版本的Upcaster类含:fromType+fromVersion → 转换后的newType+newVersion。多版本链式升级:V1→V2→V3链。

Axon Upcaster示例:

public class AccountOpenedEventUpcaster extends SingleEventUpcaster {
    @Override
    protected boolean canUpcaster(SerializableType typeHolder) {
        return typeHolder.getName().equals(AccountOpenedEvent.class.getName())
            && typeHolder.getVersion().equals("1.0");
    }
    @Override
    protected IntermediateSerializedEvent doUpcaster(IntermediateSerializedEvent event) {
        return event.withData(upgrade(event.getData(), "1.0", "2.0"));
    }
}

贪心转换(Eager):读取时实时转换原子事件(如JSON的添加/删除/重命名字段)。性能比懒转换更好,但有转发耗损。

5.2 事件演化的最佳实践

  • 永不删除事件:事件链完整性是审计追踪的核心价值
  • 进化事件的追加字段非重定义:向后兼容新字段默认值
  • 避免事件的语义变化:新事件类型替代旧类型
  • 版本字段记录:event加schemaVersion字段便于区分

6. 分布式事务与ES/CQRS的协作

事件溯源天然支持幂等命令处理跨聚合的一致性。EventStoreDB对单个聚合内原子提交严格保证,跨聚合则需Saga协调:

6.1 EventStoreDB事务

EventStoreDB保证同一Aggregate Stream上多个事件的乐观并发控制:每次Append需携带expectedVersion号。若当前stream末尾版本≠expectedVersion(意味着并发修改),抛出WrongExpectedVersionException。这就是乐观锁CAS语义。

6.2 Saga+事件溯源

跨聚合事务通过Saga协调。Saga的每个步骤是聚合上的命令,命令执行产出事件,后续步骤订阅事件触发。

以电商下单为例:

Saga-OrderCreated:
Step1 →(OrderService) OrderAggregate.Create Result Event(OrderCreated)
Step2 →(PaymentService) PaymentAggregate.Authorize PaymentOnOrderCreated Event
Step3 →(InventoryService) InventoryAggregate.ReserveStockOnPaymentReserved Event
...

若Step3库存失败,向PaymentService发送CancelPaymentCommand,PaymentAggregate产出PaymentCancelledEvent弥补。这就是Saga套件之间的事件驱动协调。

7. 设计合适的事件持有模型

事件溯源+ CQRS系统的设计仍需遵循DDD(Domain-Driven Design)原则:

  • 聚合是一致性边界:聚合内数据强一致
  • 聚合尽量小:延迟内部事件少→读放快
  • 跨聚合引用ID引用:非对象引用,避免锁膨胀
  • 聚合间一致性通过Saga:异步最终一致

7.2 事件的四个层次

Domain Events(领域事件):核心业务含义,是核心建模(如OrderPlaced)

Integration Events(集成事件):跨服务通知,承载最少信息(OrderId+Status)

Application Events(应用事件):投影和内部组件交互(如PortfolioItemAdded/Removed)

Infrastructure Events(基础设施事件):系统级别(如HealthCheckFailed、RecurringCommandSent)

7.3 事件粒度原则

小事件优势:更易演化和重放,精度更高(ItemAdded比CartCheckedOut容易组合)

小事件代价:事件数量膨胀,加载重放开销大

实战经验:电商购物车用"ItemAdded+ItemRemoved"小事件;但银行流水可用"Transacted"大事件避免1亿个读模型维护。

8. 性能优化工程实践

8.1 投影的Resume机制

EventStoreDB的$all流配合持久订阅(Persistent Subscription)可管理读取进度。subscription自带checkpoint(已处理到哪个位置),崩溃后从checkpoint继续。Axon的StreamingEventProcessor(SEP)提供类似的tracking token机制。

8.2 读写端异步流水线

事件存储 → 事件总线(Kafka/RabbitMQ) → 投影微服务(Kafka消费者组) → 读模型存储(Elasticsearch/Redis/PostgreSQL)。

关键优化点:

  • 批量写入:投影消费端按批flush,减少IO flush次数
  • 并行投影:多个投影服务彼此独立消费相同事件流,利用Kafka分区并行
  • 融错降级:投影故障时事件在消息队列中堆积,恢复后追赶(replay)

8.3 缓存与物化视图

读取端可借助Redis缓存热点聚合读模型(如用户Session数据),物化视图定期(秒级)更新。对强一致性要求高的同步(如银行余额),投影需要快照+版本号判断可见性。

9. 适用场景与陷阱

9.1 EventSourcing+CQRS的推荐场景

  • 高审计需求:金融账务、医疗记录、政府监管——需要操作全历史
  • 状态回溯/时间旅行:查询任意时间点系统状态
  • 写读极端不对称:写入极少读出极其频繁(如订单创建百万QPS读详情)
  • 复杂下流诉求:同一事件流可触发通知、分析、报表多个订阅者
  • 实时协作:协同编辑、多人任务追踪,事件流天然支持操作CRDT

9.2 不适用场景

  • 状态不重要的数据:缓存、会话等临时状态不需要历史
  • 简单的CRUD应用:传统ORM+SQL即可满足,不要过度设计
  • 写性能极高:事件存储较多的IO开销可能影响性能

9.3 典型陷阱

  • 事件结构僵化:事件模式经常变更,始终没有合适的演化策略
  • 投影跟不上写速度:热配镜的订阅速率不如事件产生速率
  • 快照设计错误:快照触发过晚导致现射回放时间过长;过晚导致存储膨胀
  • 领域建模不足:事件只是当前状态的静态快照投影,丢失业务含义
  • 强一致假象:约定"事件写入后立即可见"但实际投影有延迟

10. 总结

事件溯源+CQRS不是银弹,但它提供了三个有价值的工程能力:

  • 时间维度:事件流是系统状态变化的完整历史
  • 读写分离:双模型分别优化,互不干扰
  • 分布式协作:事件是微服务间异步协作的黏合剂

深层价值:它迫使开发者在建模时思考"哪些事情真的发生了"而非"当前状态是什么"。这种思维方式本身就是领域建模的进步。

点赞(0) 打赏

评论列表 共有 0 条评论

暂无评论
立即
投稿

微信公众账号

微信扫一扫加关注

发表
评论
返回
顶部