Kafka学习笔记(四):消息可靠性、幂等与事务
Kafka学习笔记(四):消息可靠性、幂等与事务
订单事件已经发出,消费者也能收到。接下来最关心的往往是三个问题:
- Broker 宕机后,消息还在吗?
- 消费者重试,会不会重复扣库存?
- 订单已经写入 MySQL,但 Kafka 发送失败,怎么补回来?
这三个问题发生在不同位置,需要不同的机制解决。
本篇沿用 Kafka 4.0 系列与第二篇的 Java 示例。多副本部分使用“三个 Broker、三个副本”的架构例子;第二篇的单节点环境只能验证单副本下的基本行为。
系列导航:
一、可靠性要沿着整条链路分析
先画出订单事件经过的位置:
1 | 订单写入 MySQL |
逐个位置追问“这里崩溃了会发生什么”,比直接背“Kafka 保证消息不丢”更有用。
| 环节 | 典型故障窗口 | 主要应对方式 |
|---|---|---|
| 数据库到生产者 | 订单提交了,事件还没发送 | Outbox、CDC、补偿 |
| 生产者到 Broker | 发送失败,或者结果不确定 | 检查结果、重试、稳定事件 ID |
| Broker 复制 | Leader 失效,副本未跟上 | 副本、ISR、ACK 策略 |
| 消费业务处理 | 先提交了 Offset,业务还没落库 | 对齐提交与业务完成 |
| 业务到 Offset 提交 | 落库成功,但 Offset 未提交 | 消费幂等 |
消息在 Kafka 中保存可靠,与业务最终执行正确,是两个需要一起完成的目标。
二、acks:生产者等待什么确认?
2.1 三种确认方式
| 配置 | 生产者等待什么 | 典型故障风险 |
|---|---|---|
acks=0 |
不等待 Broker 确认 | 应用无法据此知道 Broker 是否接收 |
acks=1 |
Leader 本地日志写入确认 | Leader 失效时,尚未复制的记录可能丢失 |
acks=all |
当前 ISR 所需的复制确认 | 保障程度还取决于 ISR、副本和故障范围 |
all 也可以写成 -1。官方生产者配置说明了确认语义。
这里的“写入日志”不要直接等同于“每次都完成磁盘强制刷盘”。Kafka 的可靠性依靠复制等机制,并非默认对每条消息执行 fsync。官方副本设计讨论了这一存储假设。
2.2 acks=all 等于永远等待三个副本吗?
不是。要看当前 ISR,不能只看配置的副本总数。
假设:
1 | replication.factor = 3 |
在这个经典 ISR 示例中:
| 当前 ISR | 写入结果如何理解 |
|---|---|
| A、B、C | 等待当前三个同步副本完成所需复制 |
| A、B | 满足最低要求,可在两者完成所需复制后确认 |
| A | 不满足最低要求,写入不能成功确认 |
min.insync.replicas=2 是最低门槛,不是只要任意两个副本确认就立刻返回。 配置与异常行为见官方 Broker 配置。
2.3 为什么最低门槛很重要?
如果 ISR 只剩 Leader,而且最低要求仍为 1,那么 acks=all 也可能在只有一份同步数据时成功。
提高最低同步副本数,意味着副本不足时宁可拒绝写入,也不继续降低冗余。与此同时,调用方必须能够处理写入失败,否则保护 Broker 数据的配置会变成应用侧丢弃事件的原因。
这也是为什么第二篇的单副本环境,即使配置了 acks=all,也不能用于证明节点容错。
三、发送超时:结果可能不确定
看这个过程:
1 | Producer 发送事件 E1001 |
失败结果不总能证明消息没有写入。 这时重新发送可能产生重复,直接放弃又可能漏掉确实没有到达的事件。官方投递语义分析指出了这个故障窗口。
应用应当保留事件身份,检查发送结果,并根据失败类型进行恢复。delivery.timeout.ms 可以控制一次发送在客户端中可用的整体交付时间,但它不会为业务生成补偿事件,也不会把数据库事务一起回滚。
四、生产者幂等:解决客户端重试重复
4.1 enable.idempotence 的作用
同一批记录因通信问题被客户端重试时,Broker 可以通过生产者 ID 与分区序列号识别重复写入。官方设计文档解释了这一机制。
学习示例显式配置:
1 | acks=all |
幂等开启要求 acks=all、retries>0、max.in.flight.requests.per.connection<=5。Kafka 4.0 在没有冲突配置时默认启用幂等,显式设置便于表明应用意图。生产者配置文档列出了这些约束。
max.in.flight 也不必一律改成 1;开启幂等并满足约束时,允许值范围内仍可以保持重试场景下的分区顺序。
4.2 它不能识别“同一笔业务”
考虑主动发送两次:
1 | producer.send(new ProducerRecord<>("order-events", "O1001", sameValue)); |
这是两次应用层发送,Kafka 不会因为 Key、Value 相同,就认为第二条应该删除。
| 重复来源 | 生产者幂等能否直接解决 |
|---|---|
| 同一客户端发送批次的协议级重试 | 可以处理其作用范围内的重复 |
| 用户重复点击,应用创建两次事件 | 需要应用识别业务请求 |
| 应用重启后重新发送同一业务事件 | 需要稳定事件身份与恢复策略 |
| 消费者处理后未提交 Offset,再次消费 | 需要消费业务幂等 |
所以 eventId 仍然需要保留,而且同一次业务事件的重试必须复用同一个 ID。
五、消费者幂等:把去重与业务放进同一个事务
5.1 重复消费是怎么发生的?
1 | 消费 E1001 |
这时 Kafka 的生产者幂等已经帮不上忙。问题发生在数据库写入与消费进度之间。
5.2 设计一个去重表
以下是 MySQL/InnoDB 的教学表结构:
1 | CREATE TABLE consumed_event ( |
联合主键中的 consumer_name 代表业务处理方,例如 inventory-service。
为什么不只用 event_id?同一订单事件可能需要库存、积分两个服务分别处理,库存已经执行不应该让积分误以为自己也执行过了。这个标识应当按业务语义设计,而不是每次重放随意换个名字。
5.3 正确的执行顺序
以一个需要记录消费身份的业务为例:
1 | 开启数据库事务 |
应用伪代码:
1 | try: |
这里必须精确识别去重表指定约束的重复,不能把所有数据库异常都当作“已处理”。业务更新影响行数、库存条件和状态转换也必须校验;更新失败应让整个事务回滚。
5.4 为什么事务必须把两者包住?
| 故障时机 | 去重记录与业务同事务时的结果 |
|---|---|
| 插入去重记录后,业务失败 | 一起回滚,下一次仍能尝试 |
| 事务成功后,Offset 提交失败 | 重放触发唯一约束,避免再次执行业务 |
| 两个请求同时处理同一事件 | 唯一约束协调竞争,成功结果只落一份 |
这套推导的前提是:去重与业务写入发生在同一个支持事务的数据库中,提交进度不会越过失败记录,而且去重记录保留时间覆盖重试与重放窗口。
5.5 Redis SETNX 为什么不能单独解决?
如果先执行 SETNX 记录“已处理”,再扣库存,中间崩溃后,重试可能被那个标记挡住,库存却没有扣。
反过来先扣库存再写标记,也有重复扣减的窗口。分布式锁能控制同时执行,但锁过期、释放或进程崩溃后,不能独自证明业务已经永久成功。
发短信、调用外部支付等副作用同样无法被本地数据库事务自动包住,需要下游支持幂等请求键,或额外的任务记录和对账补偿。
六、Kafka 事务:原子地提交 Kafka 内部的工作
6.1 哪类问题适合 Kafka 事务?
例如消费订单事件,计算一条风控结果,再写入另一个 Topic:
1 | 读取 order-events |
如果输出已经写入,但输入进度没提交,重启就可能再生成一条输出。Kafka 事务可以把输出记录和输入消费组进度放在同一个 Kafka 事务中提交。
6.2 一个事务处理周期
1 | 初始化事务生产者:initTransactions() |
这是流程示意,API 的初始化、异常分类和事务结束方式可参考KafkaProducer API。输入消费者应关闭自动提交,进度通过事务提交,输出读取方应配置:
1 | isolation.level=read_committed |
否则默认的读取方式可能读到随后被中止的事务记录。官方消费者配置说明了可见性规则。
事务中止并不会自动把当前消费者的内存读取位置退回去。 可恢复失败时,要中止事务并重新定位输入,或重建客户端从已提交位置继续;不能直接跳到下一批。被 fencing 等不可恢复错误拒绝的实例,需要停止并关闭客户端。
6.3 transactional.id 怎么理解?
它标识一个逻辑事务生产者。通常需要在实例重启后保持稳定,同时让并行实例使用不同 ID。相同逻辑 ID 的新实例可以使旧实例失去事务写入资格,避免旧实例继续写入。官方事务生产者说明介绍了恢复与 fencing。
单节点练习还需调整事务内部 Topic 的副本要求;第二篇未配置这些参数,因此不能直接拿流程示意在那个环境里验证完整事务。
6.4 Exactly-Once 的边界
对于 Kafka 到 Kafka 的处理,在正确的事务、输入进度和读取隔离配置下,可以实现对应边界内的 Exactly-Once 处理语义。
但是,Kafka 事务不会自动把 MySQL 更新、短信发送或第三方 HTTP 调用一起纳入原子提交。 外部副作用仍需要自己的幂等与一致性方案。官方事务设计讨论了输入进度与输出提交的关系。
七、MySQL 成功,Kafka 失败:Outbox 解决双写窗口
7.1 两种直接双写都有问题
1 | 先写数据库,再发消息 → 数据库成功后,发送可能失败 |
即使在方法上加本地数据库的 @Transactional,也不会因此让 Kafka 的确认结果与 MySQL 的事务提交变成同一个原子操作。
7.2 本地事务只做一件能保证的事
将订单记录与待发送事件放进同一个数据库事务:
1 | MySQL 本地事务 |
这是一种基于前面故障窗口分析得到的工程方案。它把“订单存在,就应该有对应事件”的约束先落在数据库中,再通过可恢复的异步投递完成后半段。
一个简化的待发送表:
1 | CREATE TABLE outbox_event ( |
7.3 投递器仍可能重复发送
1 | 发送成功 |
因此每次投递复用同一个 event_id,消费端继续做去重。Outbox 帮助恢复未完成投递,不能消除所有重复。
这个简化表结构还不是完整投递器。多实例抢占、超时重试、失败记录、清理策略都需要设计。如果同一订单要求严格顺序,投递器还要保证该订单事件的发送顺序;相同 Kafka Key 无法修正源头已经发反的顺序。
八、失败重试:不能让后续 Offset 越过错误
8.1 可恢复错误与不可恢复错误
| 失败 | 一种处理思路 |
|---|---|
| 下游暂时超时 | 有限次数重试,保持分区完成边界 |
| JSON 格式错误 | 记录错误原因,转入人工检查或错误 Topic |
| 业务数据暂未到齐 | 根据业务设计延后处理与补偿 |
| 认证、配置长期错误 | 报警并修复,避免无限空转 |
一个始终失败的事件可能挡住整个分区。如何绕过它,是业务语义选择:如果它是订单“创建”事件,后面的“支付”事件能否独立执行?
8.2 失败消息转移也有提交窗口
把失败事件发到 order-events-dlt 后,才能考虑推进原分区进度,而且要保留原 Topic、Partition、Offset、eventId 与错误原因。
直接发送错误消息、再提交输入进度,同样有双写窗口。在 Kafka 内部可以用事务把两者结合;采用普通发送时,需要确认投递成功并允许错误消息重复,同时做补偿与监控。
把消息移到错误 Topic 代表它被隔离了,不代表原业务已经完成。
对于严格顺序业务,将失败事件送进独立重试 Topic、让后续事件继续,可能打乱原分区的业务执行顺序,需要另外设计阻塞、版本校验或按实体恢复机制。
九、把可靠性变成可验证的实验
不用一开始就搭复杂集群,可以先验证消费者这一侧:
- 在测试业务完成后、提交 Offset 前人为中断进程,观察重启后的重复。
- 为事件添加稳定的 eventId,将去重与业务写入放入同一个数据库事务,再做同样的中断。
- 模拟某条记录失败,确认代码不会继续提交越过这条记录的进度。
- 保留一条待发送 Outbox,模拟投递失败,再观察恢复后能否继续发送。
多副本故障实验则需要多个 Broker:观察 ISR 缩小、Leader 切换,以及同步副本不足时生产者是否正确保留并恢复失败事件。
最后用这张职责表检查设计是否缺了一层:
| 机制 | 主要负责什么 |
|---|---|
| 副本 + ACK + 最低 ISR | Kafka 存储链路的冗余与确认条件 |
| 生产者幂等 | 协议级发送重试的重复写入 |
| 业务事务后提交 Offset | 处理失败后的重新尝试 |
| 消费业务幂等 | 重复到达时避免重复业务结果 |
| Kafka 事务 | Kafka 内输出记录与输入进度的原子提交 |
| Outbox + 恢复投递 | 数据库提交到消息发送之间的双写窗口 |
学习 Kafka 的可靠性,最有价值的练习就是选一个故障时机,把数据、事件和 Offset 各自的状态写出来,再检查恢复后会不会漏处理或重复产生业务效果。