CQRS/Event Sourcing模式下如何实现视图表持续更新
CQRS/Event Sourcing模式下视图表同步落地方案
首先直接给结论:日常同步流程里完全不需要每次执行命令就读取全量事件重算视图表,这是很多刚接触ES模式的开发者最容易踩的认知误区。KSQL只是流式投影的其中一种实现,完全可以基于普通关系型数据库搭出稳定、高性能的同步链路,核心逻辑围绕「增量投影+进度断点」展开,和用不用大数据组件没有必然关系。
核心同步逻辑(关系型数据库可直接落地)
- 先给事件存储加全局有序标记
不管是用专用事件存储,还是直接用MySQL/PostgreSQL这类关系库存领域事件,必须给每个事件分配一个全局单调递增的偏移量(自增主键、全局序列号都可以),不能仅靠业务聚合ID排序,避免跨业务实体的事件处理乱序。
关系库存事件的通用表结构参考:CREATE TABLE domain_events ( global_offset BIGINT AUTO_INCREMENT PRIMARY KEY, event_id VARCHAR(64) UNIQUE NOT NULL, aggregate_id VARCHAR(64) NOT NULL, aggregate_type VARCHAR(32) NOT NULL, event_type VARCHAR(64) NOT NULL, event_data JSON NOT NULL, occurred_at DATETIME NOT NULL, INDEX idx_aggregate_query (aggregate_type, aggregate_id) ); - 给每个读模型单独存消费断点
每个视图表(读模型)对应一个独立的投影任务,每个任务单独维护自己的消费进度:也就是记录「当前这个投影已经处理完哪个全局偏移量之前的所有事件」,进度存在单独的断点表中即可,表结构很简单,存投影名称、最后处理的偏移量、最后更新时间三个字段就够。
举个例子,面向C端的用户订单列表视图表v_user_order_view,对应的投影任务就在断点表里存一条记录:投影名user_order_list,last_processed_offset=23679,代表偏移量1到23679的所有相关事件都已经同步到视图表里了。 - 增量拉取+幂等更新
投影任务每次启动/轮询的时候,先读自己的消费断点,只拉取偏移量大于断点值的新事件,按全局顺序逐批处理,每处理完一批事件就更新对应的断点值。
这里必须做幂等兜底:如果出现事件处理完、断点还没更新服务就宕机的情况,下次重启会重复拉取这部分事件,只要视图表更新逻辑做了幂等(比如用事件ID做唯一约束,重复处理同一事件不会重复写入/更新错误数据),就不会出现一致性问题。
中小流量场景下根本不需要引入消息队列、流处理引擎,直接用定时任务轮询事件表拉取增量就行,轮询间隔设100ms到1s,延迟完全能满足业务需求。如果事件量特别大,可以给投影任务做分片,比如按聚合ID哈希拆成多个消费实例并行处理,只要保证同一个聚合ID的事件始终路由到同一个实例,就不会出现处理乱序。 - 快照优化解决单聚合事件堆积问题
针对单条业务实体(聚合根)事件量持续增长的问题,不用全量扫该实体的所有事件:每隔固定数量的事件(比如每50个事件)就存一份该聚合的最新状态快照,后续写侧重建聚合状态做业务校验的时候,直接加载最新版本的快照,再追加处理快照版本之后的增量事件即可,不需要从头遍历该聚合的全部历史事件。注意快照是优化写侧聚合加载性能的,和读模型投影流程无关,不要混用。
常见避坑点
- 不要把读模型更新放到命令处理链路里:命令处理属于写侧逻辑,只负责校验业务规则、生成并持久化领域事件,读模型更新是完全异步的旁路流程,和写链路解耦才能体现CQRS读写分离的性能优势。
- 不要搞事件和读模型双写:绝对不能在写业务数据的时候同步更新视图表,只要其中一边写失败就会导致读写数据不一致,所有读模型的更新必须以事件存储中已经持久化成功的事件为唯一输入。
- 全量重建是常规操作,不用刻意回避:如果投影逻辑有调整需要重建读模型,直接清空对应视图表、把对应投影的断点重置为0,从头跑一遍增量消费逻辑即可。只要做了批处理+分片,哪怕是几亿规模的事件,重建速度也很快,日常运行永远走增量消费,不会出现全表扫描的性能问题。

内容的提问来源于stack exchange,提问作者juanmorschrott
相关产品推荐
相关产品推荐

