事件溯源/CQRS中多实例读模型生产者的K8s部署问题与应对
读模型生产者多实例重复消费与单点故障解决方案
一、搞定多实例重复处理事件:实现分布式消费协调
1. 利用Event Store原生的消费组/竞争订阅
大部分成熟的Event Store(比如EventStoreDB)都支持竞争消费者模式的持久化订阅:把所有ReadModelProducerService实例挂到同一个订阅下,Event Store会自动把事件分发给可用的实例,保证每个事件只被处理一次。启动时记得配置订阅从上次检查点开始消费,别每次都重放全量历史,不然启动成本太高。
如果用的是类似Kafka的事件存储,直接用消费组(Consumer Group)就行,组内实例自动分片消费,天然避免重复。
2. 自定义进度追踪+分布式锁兜底
如果你的Event Store不支持原生协调,自己搭一套也不难:
- 在ReadModel数据库里加一张
event_consumption_progress表,记录每个读模型对应的最后处理事件ID/位置 - 实例启动先读这张表,从上次的位置继续消费
- 处理每个事件前,用分布式锁(比如Redis Redlock,或者K8s Lease)抢处理权,抢到的才处理,处理完立刻更新进度表。这样就算多实例同时跑,也不会重复处理同一件事。
二、解决单实例单点故障:高可用优化
1. 给K8s加精准的健康检查
K8s默认的存活探针只会看进程在不在,管不了服务“假死”(进程活着但不处理事件)。你得自定义探针:
- 存活探针:暴露一个接口,返回最近一次处理事件的时间戳,如果超过5分钟没处理新事件,就返回失败,K8s会自动重启这个实例
- 就绪探针:检查服务是否已经追平Event Store的事件进度,只有追平的实例才允许加入消费队列,避免刚启动的实例乱抢事件
2. 多实例做主节点选举
搞个主备模式,同一时间只有一个实例干活,其他实例待命:
- 用K8s的Lease API做选举:每个实例尝试抢一个命名的Lease,抢到的就是主节点,负责消费事件;主节点每隔几十秒续租一次Lease,要是它假死了,续租超时,备用节点就会自动抢Lease上位,继续处理事件,无缝衔接。
3. 架构拆分解耦
别把所有读模型的生产逻辑都塞在一个ReadModelProducerService里:
- 每个读模型单独做一个小的消费服务,这样一个服务故障,只影响对应的读模型,不会整个应用瘫痪
- 配上HPA(水平Pod自动扩缩容),根据事件堆积数自动扩容,故障时K8s会快速拉起新实例替换
三、实操建议
- 优先用Event Store原生的消费协调机制,自己写代码容易踩坑
- 一定要测故障场景:手动把主节点的进程挂起(不是杀掉),看K8s能不能检测到假死,备用节点能不能自动接管
- 消费进度一定要持久化,别存在内存里,不然实例重启就得重放全量历史
内容的提问来源于stack exchange,提问作者Arczewski
相关产品推荐
相关产品推荐

