幂等写入路径:在消息重复时保持集成管道正确性
一笔对不平的账
先讲一个真实案例。一次大促结束后对账时,财务发现一批订单状态不对:我们这边显示“已支付”,但财务系统并没有入账。
顺着链路追查后我们发现,支付网关在高峰期重试了一次,同一条支付通知被推送了两次。第一次被正常处理;第二次到达时,消费者刚好处于滚动重启过程中,消息被重新入队并再次被消费。
中间状态表被写了两遍。下游做了自己的去重所以没受影响,但我们这边的对账报表逻辑被多出来的一条脏数据搞崩了。
这件事本身不算严重,但它让我们清楚地认识到:在分布式环境下,消息就是会重复。网络重传、队列重投、消费者重启、上游超时重发——你根本阻止不了。
唯一能做的,就是让系统在收到一次消息和收到五次消息时,行为保持一致。这就是幂等性。
幂等键选什么,才是关键
很多人以为幂等就是“加个唯一键”。对,也不对。关键在于你究竟以什么为键。
一开始我们直接拿消息队列的消息 ID 作为幂等键。看起来合理:每条消息不都有 ID 吗?
但坑就在这里。如果上游手动重推一次,或者它自己的重试机制触发了,就会生成一个全新的消息 ID。在消息层面看,这像是两条“不同的消息”,但在业务层面它们是同一件事,所以去重直接失效。我们在这个坑里栽了两次,才彻底修好。
后来我们改成要求所有上游在事件体里带上业务级别的幂等键。幂等键怎么构造,不同场景不一样,但核心原则只有一个:必须能唯一标识“这个业务事件发生过一次”。举几个例子:
注意,幂等键往往是复合的。同一个订单会触发多个状态变更,所以不能只用订单 ID:“已支付”和“已发货”是两件不同的事,必须加上状态或版本号来区分。
“唯一标识一个实体”和“唯一标识一次事件发生过”是两码事,很多人会搞混。
推行这个标准也不是一帆风顺。有些上游团队会觉得:“凭什么要我改数据格式来配合你的管道?”遇到这种情况,我们只能请架构师出面对齐,解释为什么这件事必须在源头做。
如果你想在中间层打补丁,那是补不完的,因为你根本不知道上游的业务语义。
去重检查到底放在哪
一开始我们在业务代码里做去重:处理消息前先查一下这个键是否已经处理过。功能上没毛病,但存在竞态条件。
当消费者的处理逻辑比较复杂时(一个事件要写三张表、调用两个下游服务),从去重检查到真正写入之间有一段处理时间。如果重复消息恰好在这个窗口内到达,就会出现“检查时没有重复,写入时却有重复”。在高并发下,这种情况比你想象的更常见。
后来我们把去重下沉到数据库层,用唯一约束兜底。先建一张去重表:
[LOADING...]
处理逻辑大致如下:去重记录和业务数据在同一个事务里写入,要么都成功,要么都回滚:
这样,去重和业务写入就是原子的,不再有竞态窗口。
顺便说一句,这张去重表会一直增长,所以需要定期清理。我们搞了一个定时任务,每天晚上清理 30 天前的记录。30 天这个数字是根据我们业务中实际的消息重复窗口定的:绝大多数重复投递都发生在几分钟内,30 天绰绰有余。
如果业务写入跨多个数据源(写完数据库还要调外部 API),单靠数据库事务就覆盖不了,需要补偿机制。这部分后面讲韧性的时候再展开。
重复消息来了:忽略、覆盖还是合并
怎么处理取决于业务语义,没有标准答案。实践中我们遇到过三种策略。
忽略用得最多:检测到重复就直接丢掉,什么都不做。它适合天然幂等的操作,比如“把用户状态设为已激活”,做一次和做十次效果一样。在我们系统里大约 70% 的事件类型走这条路。它实现最简单,也最不容易出 bug。
覆盖适合最终状态同步,比如从 CRM 同步客户的最新联系方式。但这里有个坑:如果事件乱序到达(先到的那条反而是新数据,后到的是旧数据),盲目覆盖会把新数据回滚掉。这类 bug 排查起来特别恶心,因为数据“看起来是正常的”,只是不是最新的,你可能很久都不会察觉。所以我们给覆盖加了一个版本号检查:
这本质上是一个简化的 Last-Write-Wins:只有更高版本才能覆盖。这个版本号必须由源系统生成。你不能在中间层自己造,因为中间层并不知道源系统在什么时候更新了什么。
合并是最复杂的:把旧数据和新数据合到一起。它适合增量数据(给订单追加备注、往购物车加商品这类场景)。它的实现复杂度最高,也最容易出 bug。就拿“追加备注”来说:如果同一条备注因为重复投递被追加了两次怎么办?那你就得在合并逻辑里再做一层去重。套娃,基本就是这样。所以能不用就不用。
跨系统把身份串起来
做集成的朋友大概都遇到过:同一个实体,不同系统有不同的 ID。一个客户在 CRM 里叫 CRM-00123,到 ERP 里变成 ERP-C456,到 WMS 里又是 WMS-CUST-789,我们平台上还有一套自己的内部 ID。
当事件从 CRM 流向 ERP 时,你得知道 CRM-00123 对应 ERP 里的哪条记录。我们在集成平台里维护了一张 ID 映射表,用平台内部的 internal_id 把各个系统的 external_id 串起来:
[LOADING...]
事件进入管道后做的第一件事,就是查映射表拿到目标系统的 ID。如果查不到(比如创建事件,下游还没有对应记录),就标记为“待创建”,等下游创建完成后,我们再把映射关系写回去。
这里有个细节:映射回写本身也可能失败或重复,所以对这表的写入也得是幂等的。我们用 INSERT ... ON CONFLICT DO NOTHING。
另一个非常重要的东西是关联 ID(correlation ID)。每个业务流程(比如“创建一个订单并同步到所有系统”)都会有一个全局唯一 ID,贯穿该流程的每一个事件和每一行日志。
排查问题时,一个 grep 就能把散落在几十个系统、几百行日志里的信息串起来。用过的人都知道,出故障时这东西有多重要。没有它,你只能瞪着几十个系统各自独立的日志,完全不知道谁和谁有关联,定位一个问题从半小时变成半天。
怎么证明你没丢数据
每隔几天老板就会问同一个问题:“你怎么证明管道没有丢数据?”分布式系统确实没法从数学上证明,但你可以用工程手段把置信度推得很高。
我们做了端到端的事件追踪。每个事件从进入管道到离开,每经过一个处理节点都会留一条追踪日志:
这些追踪日志写入 Kafka,并长期归档到对象存储。在此基础上,我们跑两种对账。
实时对账每五分钟扫描一次,找出“已接收但未完成”的事件。如果某个事件在接收后超过十分钟还没有 COMPLETED 记录,就触发告警。这能很快暴露“卡住”的事件,比如消费者死锁,或者下游接口一直不返回。
离线对账每天晚上做一次全量比对。拿上游声称已发送的事件列表,和我们实际处理完成的事件列表做 diff,找出“上游发了但我们没处理完”以及“我们处理了但上游根本没发”。前者意味着丢数据;后者可能是重复或幽灵记录。两种对账互补:实时对账快但覆盖有限,离线对账慢但能兜住边界情况。
有了这套追踪系统,对老板的答复就不再是“不可能丢数据”,而是“如果丢了一条,我们五分钟内发现,十分钟内定位,半小时内恢复”。
关于作者
Yuelin Ou 是一名数据与 AI 工程师,工作重点是幂等写入路径、分布式管道韧性,以及在不破坏正确性保证的前提下扩展企业集成系统。她拥有罗切斯特大学数学学士学位,辅修计算机科学。个人网站:yuelinou.com。
本文《幂等写入路径:当消息重复时如何保证集成管道正确》首发于 foojay。