CQRS分离存储架构下的数据同步逻辑与实现疑问
CQRS架构下写端(事件存储)与读端(Elasticsearch)的同步方案
一、核心同步流程
基于你给出的架构,同步链路是固定的:
- 写端产生业务事件后,事件存储将事件流式推送到消息中间件
- 读端的专属同步组件监听消息中间件的事件流,负责把事件转换为读模型,同步更新Elasticsearch
二、是否需要独立组件?
必须写。这个组件就是CQRS里的投影(Projection),是连接写端事件流和读端存储的核心。事件存储只负责记录事实,不会主动维护读端数据,必须靠专门的组件做事件到读模型的转换、更新。
三、单事件下的状态同步逻辑
如果当前只有单个事件,处理逻辑很直接:
- 投影组件消费到事件后,根据事件携带的实体ID,在Elasticsearch中查找对应文档(首次事件则文档不存在)
- 直接用事件数据构建或更新读模型:比如事件是
UserCreated {id: 1, name: "Alice"},就创建ID为1的用户文档;如果是UserNameUpdated {id:1, newName: "Alicia"},就找到ID1的文档更新name字段 - 将更新后的文档写入Elasticsearch完成同步
四、最终状态的来源与获取逻辑
- 最终状态是什么:是实体从创建到当前所有事件的累加结果。比如用户最终状态=创建事件+改名事件+手机号更新事件的合并值。
- 怎么获取:
- 常规场景:投影组件消费完整的事件流,从实体的第一个事件开始,逐个将事件的变更应用到读模型上,最终叠加出当前的最终状态。
- 单事件场景:如果只有一个事件,那这个事件本身就是最终状态的全部内容。
- 关键逻辑:不是靠单一事件传递最终状态,而是靠事件流的顺序叠加构建——每个事件只记录一次变更,投影组件把这些变更依次应用,得到最新状态。
五、快照机制的性能优化
会用快照优化性能,尤其是当实体事件量极大时:
- 当某个实体的事件积累到阈值(比如1000个),投影组件可以在处理完一批事件后,将当前的读模型状态保存为快照
- 后续如果需要重新构建该实体的读模型,不需要从头消费所有事件,只需加载最新快照,再消费快照之后的事件即可,大幅减少计算和IO开销
- 快照通常存在Elasticsearch本身,或单独的对象存储里,也可以作为特殊事件类型存在事件存储中
内容的提问来源于stack exchange,提问作者Arabaaa
相关产品推荐
相关产品推荐

