系列目录
- MySQL 索引与慢查询:B+ 树如何减少扫描
- 声明式事务之下:InnoDB 的 MVCC 与锁
- 接入层:Nginx 反向代理与 OpenResty 的边界
- Web 容器与 Netty:线程模型之下的 IO 模型
- Redis(上):缓存用法与单线程模型的限制
- Redis(下):超出缓存用途的用法:锁、队列与排行榜
- 数据库访问层:连接池与 MyBatis 的显式 SQL
- Kafka(上):吞吐的来源是顺序 IO
- Kafka(下):生产端、broker 与消费端的可靠性配置
- RPC 框架:像本地调用一样调远程的代价
- 消息语义:按 at-least-once 设计业务代码(本篇)
使用现场:三个业务里的重复消息
(Kafka 异步化与削峰)提到消费幂等约定:去重表加唯一消息 ID、定时对账。那篇讲的是云盘两个链路的业务接入方式。本文比较 Kafka、RocketMQ、RabbitMQ 的语义保证,并说明业务代码应采用的假设;引用内容不重复展开。
重复消息在生产里反复出现过。下面三个案例形态不同,均由至少一次投递与消费端幂等处理未覆盖共同导致。案例按记忆重构,第一个有真实原型,后两个脱敏改写。
重复通知。 上传完成后的多端通知链路,消费端收到消息后调用各客户端推送通道。某次下游推送通道响应变慢,消费端重试超时,消息被重新投递。推送通道自身没有按消息 ID 去重,用户在手机上收到了两次「上传完成」通知。去重表只能避免消费侧重复写入业务数据,不能避免外部通道重复推送。本地事务的去重范围不包含外部副作用,外部接口也要按消息 ID 实现幂等。2016 第 6 篇代码注释中「通知接口自身仍要能按消息 ID 去重」的要求未被落实,导致了这次重复通知。
重复计数。 一个统计消费组把点击事件累加写入一张计数表,逻辑是 UPDATE counter SET cnt = cnt + 1 WHERE dim = ?。这条 SQL 本身不幂等:执行一次和执行两次结果不同。消费端崩溃重启后,offset 还停在处理前,这批消息被重新拉取重新累加,当天某个维度的计数比实际偏高。对账发现生产侧原始事件数与计数表对不上,差额对应重启窗口。修复方式是把累加改成按事件 ID 的去重表加 INSERT IGNORE,计数由去重记录聚合得出,单条事件无论重投几次只计一次。
状态错乱。 一个订单状态机消费组处理「支付成功」事件,把订单从「待支付」推到「已支付」。重投发生时,消费逻辑没有校验当前状态,直接执行 UPDATE order SET status = '已支付' WHERE id = ?。对同一订单重复设置「已支付」状态时,这条 SQL 的结果相同;但同一订单还有一条延迟到达的「取消」消息(用户在支付边缘时刻点了取消)。两条消息跨分区无序,「取消」先到把状态改成「已取消」,重投的「支付成功」后到又改回「已支付」。状态机没有按版本号裁决新旧,晚到的旧事件覆盖了新状态。
三个案例都使用至少一次投递,却缺少相应的消费侧约束。通知通道没有幂等处理,计数使用了非幂等 SQL,状态机没有版本裁决;重复投递因此分别造成重复通知、计数偏高和状态回退。
端到端恰好一次难在哪里
至少一次(at-least-once)是消息系统的默认投递语义:消息至少被投递一次,可能重复。恰好一次(exactly-once)要求消息被处理且仅处理一次。实现恰好一次需要处理多个环节产生的重复。
投递链路分三段,每段都有失败窗口:
- 生产者段(窗口 A)。 生产者发消息后等待 broker 确认。网络超时或 ACK 丢失时,生产者无法判断消息是否已经写入。重试可能产生重复:broker 已经收到消息,但 ACK 未返回;不重试则可能丢失消息。生产者仅凭自身状态无法区分这两种情况。
- broker 段(窗口 B)。 Kafka(下):生产端、broker 与消费端的可靠性配置讲过,leader 确认后、follower 同步前宕机,或 unclean 选举让非 ISR 副本当 leader,已确认的消息会丢。在配置正确的前提下,broker 段可避免消息丢失;任何故障下都不丢不重不属于它的保证范围。
- 消费者段(窗口 C)。 消费者处理完消息、提交位移前崩溃,重启后重新拉取这批消息,导致重复处理。这是常见的重复来源。Kafka(下):生产端、broker 与消费端的可靠性配置把位移提交时机拆成了先提交后处理、先处理后提交、原子化三种,多数场景选「先处理后提交」,相应地会产生重复。
三段任一失败都需要协调才能达成恰好一次。两阶段提交(2PC)是经典的协调方案:准备阶段让所有参与者锁住资源,提交阶段统一提交或回滚。它的代价是协调者成为单点,准备阶段持有的锁阻塞所有参与者,吞吐随参与者数量下降。消息系统每秒处理成千上万条消息,2PC 的协调开销不适合分摊到每条消息上。分布式事务:2PC 的代价与业务补偿模式会专门讲分布式事务;在消息逐条投递的场景,2PC 的代价过高。
Kafka 0.11 提供了两项能力,但适用范围都有限:
幂等生产者(enable.idempotence=true)处理窗口 A 的重试重复。开启后,生产者拿到一个 PID(producer ID),每条消息带一个递增的 sequence number。broker 端按 <PID, 分区, sequence number> 去重,同一个生产者在同一分区内因重试产生的重复会被 broker 丢弃。限制在于:PID 是单会话的,生产者重启后拿到新 PID,旧 PID 写入的消息对新生产者来说无法去重;去重只在一个分区内有效,即「单会话、单分区」。跨会话、跨分区的重复不在覆盖范围内。
事务(transactional.id)保证「消费、处理、生产」流程内的原子性:一个生产者写多个分区时,消息要么全部可见,要么全部不可见;配合隔离级别,消费者读不到未提交的事务消息。典型场景是「从输入 topic 消费、处理、写到输出 topic」的流处理管道,输入位移提交和输出写入在同一事务里原子完成。限制在于:事务只覆盖 Kafka 内部的位移和分区写入,不覆盖下游对 DB、Redis 等外部系统的写入。消费侧把处理结果写进 MySQL 时,该写入不在 Kafka 事务内,仍需业务幂等处理。
业务代码应以消息至少投递一次、消费必须幂等为默认假设。生产端的幂等生产者和事务能减少重复来源,但不能替代消费幂等。
消费幂等的四种实现
消费幂等的目标是:同一条消息处理一次和处理 N 次,业务结果一致。四种实现各有成本和适用场景。
业务唯一键加唯一约束。 消息体携带业务侧的唯一标识(订单号、事件 ID、文件 ID 加操作类型),消费时把这个标识写进业务表,并加唯一约束。重复写入触发唯一键冲突,消费端捕获冲突直接跳过。这是最直接的做法,前提是业务表本身能容纳这个唯一键。2016 第 6 篇的去重表是它的变体:单独建一张 (msg_id, consumer_name) 表做唯一约束,业务表不改动。去重表和业务写入在同一个本地事务里提交,保证「去重记录写了但业务没写」不会单独发生。
@Transactional
public void onMessage(String msg) {
Event event = parse(msg);
// 唯一约束:(msg_id, consumer) 主键冲突时 insertIgnore 返回 0
int rows = consumedDao.insertIgnore(event.getMsgId(), CONSUMER);
if (rows == 0) return; // 已处理过,跳过
doConsume(event); // 业务处理与去重记录同事务
// 业务异常则整个事务回滚,去重记录一并回滚,重投可重新处理
}这段代码的事务只覆盖去重表和 doConsume 所在的同一个数据库。doConsume 调用外部接口(推送通道、第三方 API)时,事务回滚无法撤销已经发出的外部副作用,外部接口仍要实现幂等。
状态机校验。 适合有明确状态迁移的业务。消费前先查当前状态,判断这次事件是否合法:订单已在「已取消」状态,迟到的「支付成功」事件直接跳过。这要求业务状态本身能表达「这条消息是否已处理或已过期」,不需要额外去重表。需要维护状态迁移规则;规则错误时,幂等处理会失效。案例三的状态错乱由状态机未按版本号裁决新旧导致。修复方式是消息带单调递增的版本号,消费时比较版本号,晚到的旧版本直接丢弃。
去重表。 业务表无法添加唯一约束,或者去重维度和业务主键不一致时,可以单独建一张去重表。和「业务唯一键」的区别在于去重表独立存放:业务表存业务数据,去重表存 (msg_id, consumer, 处理时间)。这会增加一张表和一次写入,但业务表无需为幂等调整。去重表会增长,需要定期清理已处理且超过保留期的记录。去重表和业务写入必须写入同一个库;跨库时,需要处理分布式事务问题。
外部存储原子操作。 消费逻辑写 Redis 等外部存储时,用 SETNX(set if not exists)或 Lua 脚本把「判定是否已处理」和「写入」做成原子的。消息 ID 作 key,SETNX msg_id 1 返回 1 表示首次处理,返回 0 表示已处理过。成本是引入一个外部依赖,且 Redis 与业务 DB 之间不是原子的:Redis 标记成功但 DB 写入失败,下次重投时 Redis 已标记会跳过,业务实际没有处理。Redis 标记与 DB 写入不一致时仍需通过对账发现。适用场景是轻量级去重,且业务可接受通过后续对账发现这类不一致。
四种实现的复杂度和适用范围不同。业务唯一键的改造较少,前提是业务表能加约束;外部存储原子操作涉及跨存储的一致性问题。多数业务使用去重表即可满足幂等要求。
顺序性的成本
分区或队列内有序的实现成本较低;全局有序会牺牲吞吐。
- Kafka 分区内有序。 同一分区内消息按写入顺序排列,消费者按顺序消费。跨分区没有全局顺序。需要按业务对象保序时,用稳定分区键(如订单 ID、文件 ID)把同一对象的消息导到同一分区,这和 2016 第 6 篇、Kafka(上):吞吐的来源是顺序 IO的约定一致。
- RocketMQ 队列内有序。 RocketMQ 的队列(MessageQueue)对应 Kafka 的分区概念,同队列内有序。生产端按业务 key 选队列,消费端按队列串行消费。RocketMQ 4.x 的顺序消费模式(MessageListenerOrderly)锁住队列保证消费顺序,代价是消费并发受队列数限制。
- RabbitMQ 单队列有序。 单个队列 FIFO,消费者单线程消费时有序。多个消费者分担同一队列时,消息分发后不保证处理完成顺序,要严格有序只能单消费者。
三者的共同点是有序只在分区或队列内成立,超出该范围则无序。要求全局有序意味着只能用一个分区或一个队列,消费并发降为一。单分区或单队列限制了消费并发,调参无法消除这一限制。
遇到「全局有序」需求时,应先确认具体约束。许多场景只要求「同一业务对象的事件有序」,可用分区键处理;全局有序则需要接受单分区带来的吞吐损失。案例三的订单状态错乱由同一订单的两条消息被分到不同分区且未携带版本号导致。使用订单 ID 作为分区键并按版本号裁决,可以保持同一订单的事件顺序,同时避免全局有序的吞吐限制。
重试与死信:无法消费的消息怎么办
消费失败有两种:瞬时失败(下游短暂不可用、网络超时)和永久失败(消息格式错误、业务逻辑无法处理)。两者的处理方式不同。
瞬时失败通过重试处理。 重试要分级:第一次失败立即重试,后续间隔递增(指数退避),避免下游恢复期间承受更多重试流量。RPC 框架:像本地调用一样调远程的代价讲 RPC 重试放大时提过,重试次数是乘法关系,分层叠加会放大流量。消息消费的重试限制在消费侧单层,不与下游 RPC 重试叠加。重试次数有上限,3 次是常见值。
永久失败通过死信队列隔离。 达到重试上限仍失败的消息,转入死信队列(dead letter queue)。死信队列是一条独立的队列或 topic,存放无法消费的消息及其失败上下文(失败原因、重试次数、原始消息体)。主链路可继续提交位移,不会被一条持续失败的消息阻塞。一条格式错误的消息反复重试会导致后续消息全部积压。补偿任务或人工处理死信消息:修正消费逻辑后重新投递,或确认是脏数据后丢弃。
对账检查整条链路的差异。 前两类措施处理单条消息能否消费成功;对账检查生产和消费是否一致。2016 第 6 篇的对账是「事件表待发送超阈值告警 + 生产消费两侧计数差异告警」,本文给出跨 MQ 的通用模式:比较生产侧已确认消息数与消费侧已处理消息数,差额对应漏处理或重复处理;检查事件表中长期停留的记录,定位处理停在哪个环节。对账不提供实时一致性,发现问题的延迟由对账频率决定(天级或小时级)。前两类措施无法证明所有下游已经对齐;对账通过事后核对发现这种差异。
要达到「业务等效的恰好一次」,需要同时处理重复、无法消费的消息和生产消费未对齐的问题。幂等避免重复生效,死信队列隔离持续失败的消息,对账发现漏处理和未对齐。每项措施覆盖的失败类型不同,不能替代另两项。
边界清单
| 业务需求 | 应使用的模式 | 不成立的假设 |
|---|---|---|
| 允许重复投递、不允许重复生效;单分区保序 | 至少一次 + 消费幂等(去重表/唯一键) | broker 保证恰好一次 |
| 本地事务与发消息需原子 | 本地事件表 + 扫描投递 + 对账(见 RocketMQ 的事务消息与延迟消息) | 跨 DB 与 MQ 原子提交(2PC 代价无法摊到每条消息) |
| 同业务对象事件需保序 | 稳定分区键 + 消息版本号裁决 | 全局有序且高吞吐(全局有序 = 单分区 = 并发降为一) |
| 消费失败需隔离 | 重试分级 + 死信队列 + 补偿任务 | 失败消息会自动恢复 |
| 全链路最终一致 | 幂等 + 死信 + 对账三层叠加 | 任一层单独足够 |
还需注意以下边界:
- Kafka 幂等生产者只覆盖单会话单分区。 生产者重启后 PID 重置,跨会话重复不保证;跨分区重复不保证。业务侧仍需消费幂等。
- Kafka 事务不覆盖外部系统。 事务只原子化 Kafka 内部的位移提交和分区写入,下游写 DB、调外部接口仍靠业务幂等。
- 去重表与业务写入须同库同事务。 跨库时,去重记录和业务写入不在一个本地事务里,可能出现「去重写了业务没写」的不一致,需要处理分布式事务问题。
- 外部接口的幂等由接口实现。 消费端事务回滚无法撤回已发出的外部副作用,推送通道、第三方 API 要按请求 ID 去重。案例一的重复通知没有覆盖这一要求。
- 对账发现问题的延迟由频率决定。 天级对账意味着最坏一天后才发现差异,实时性要求高的场景不能只靠对账。
当时的局限
当时我只从文档了解流处理框架的恰好一次语义。Flink 的 checkpoint 机制通过屏障(barrier)对齐和状态快照实现 exactly-once,这套机制在流处理框架内部有效。但我没有在生产中使用过 Flink,也没有验证 checkpoint 的对齐协议、状态后端选择和故障恢复时的状态重建。
参考资料
- Kafka 官方文档(producer / consumer 投递语义、幂等生产者与事务章节,0.11 / 1.x 版本)
- RocketMQ 官方文档(消息顺序、消费重试与死信章节,4.x 版本)
- Kafka 异步化与削峰(2016 系列第 6 篇,消费幂等约定,本篇引用不复述)
- Kafka(上):吞吐的来源是顺序 IO(分区作为顺序性边界)
- Kafka(下):生产端、broker 与消费端的可靠性配置(三环节配置与位移提交时机)
- RocketMQ 的事务消息与延迟消息(事务消息实现与消费幂等衔接)
- 分布式事务:2PC 的代价与业务补偿模式(2PC 代价与消息最终一致)
