Kafka学习笔记(四):消息可靠性、幂等与事务

订单事件已经发出,消费者也能收到。接下来最关心的往往是三个问题:

  1. Broker 宕机后,消息还在吗?
  2. 消费者重试,会不会重复扣库存?
  3. 订单已经写入 MySQL,但 Kafka 发送失败,怎么补回来?

这三个问题发生在不同位置,需要不同的机制解决。

本篇沿用 Kafka 4.0 系列与第二篇的 Java 示例。多副本部分使用“三个 Broker、三个副本”的架构例子;第二篇的单节点环境只能验证单副本下的基本行为。

系列导航:

  1. 核心概念与架构
  2. Docker 与 Java 入门实战
  3. 消费组、Offset 与 Rebalance
  4. 消息可靠性、幂等与事务
  5. 面试高频问题总结

一、可靠性要沿着整条链路分析

先画出订单事件经过的位置:

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
订单写入 MySQL
↓
应用构造事件
↓
Producer 缓冲、发送、等待确认
↓
分区 Leader 追加日志
↓
Follower 复制
↓
Consumer 拉取
↓
业务事务提交
↓
消费 Offset 提交

逐个位置追问“这里崩溃了会发生什么”,比直接背“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
2
3
replication.factor = 3
min.insync.replicas = 2
acks = all

在这个经典 ISR 示例中:

当前 ISR 写入结果如何理解
A、B、C 等待当前三个同步副本完成所需复制
A、B 满足最低要求,可在两者完成所需复制后确认
A 不满足最低要求,写入不能成功确认

min.insync.replicas=2 是最低门槛,不是只要任意两个副本确认就立刻返回。 配置与异常行为见官方 Broker 配置。

2.3 为什么最低门槛很重要?

如果 ISR 只剩 Leader,而且最低要求仍为 1,那么 acks=all 也可能在只有一份同步数据时成功。

提高最低同步副本数,意味着副本不足时宁可拒绝写入,也不继续降低冗余。与此同时,调用方必须能够处理写入失败,否则保护 Broker 数据的配置会变成应用侧丢弃事件的原因。

这也是为什么第二篇的单副本环境,即使配置了 acks=all,也不能用于证明节点容错。


三、发送超时:结果可能不确定

看这个过程:

1
2
3
4
5
6
7
Producer 发送事件 E1001
↓
Broker 已经接收
↓
确认响应因网络故障丢失
↓
Producer 观察到超时或失败

失败结果不总能证明消息没有写入。 这时重新发送可能产生重复,直接放弃又可能漏掉确实没有到达的事件。官方投递语义分析指出了这个故障窗口。

应用应当保留事件身份,检查发送结果,并根据失败类型进行恢复。delivery.timeout.ms 可以控制一次发送在客户端中可用的整体交付时间,但它不会为业务生成补偿事件,也不会把数据库事务一起回滚。


四、生产者幂等:解决客户端重试重复

4.1 enable.idempotence 的作用

同一批记录因通信问题被客户端重试时,Broker 可以通过生产者 ID 与分区序列号识别重复写入。官方设计文档解释了这一机制。

学习示例显式配置:

1
2
3
acks=all
enable.idempotence=true
max.in.flight.requests.per.connection=5

幂等开启要求 acks=all、retries>0、max.in.flight.requests.per.connection<=5。Kafka 4.0 在没有冲突配置时默认启用幂等,显式设置便于表明应用意图。生产者配置文档列出了这些约束。

max.in.flight 也不必一律改成 1;开启幂等并满足约束时,允许值范围内仍可以保持重试场景下的分区顺序。

4.2 它不能识别“同一笔业务”

考虑主动发送两次:

1
2
producer.send(new ProducerRecord<>("order-events", "O1001", sameValue));
producer.send(new ProducerRecord<>("order-events", "O1001", sameValue));

这是两次应用层发送,Kafka 不会因为 Key、Value 相同,就认为第二条应该删除。

重复来源 生产者幂等能否直接解决
同一客户端发送批次的协议级重试 可以处理其作用范围内的重复
用户重复点击,应用创建两次事件 需要应用识别业务请求
应用重启后重新发送同一业务事件 需要稳定事件身份与恢复策略
消费者处理后未提交 Offset,再次消费 需要消费业务幂等

所以 eventId 仍然需要保留,而且同一次业务事件的重试必须复用同一个 ID。


五、消费者幂等:把去重与业务放进同一个事务

5.1 重复消费是怎么发生的?

1
2
3
4
5
6
7
消费 E1001
↓
库存扣减成功,数据库事务已提交
↓
Offset 提交之前进程崩溃
↓
重启后再次消费 E1001

这时 Kafka 的生产者幂等已经帮不上忙。问题发生在数据库写入与消费进度之间。

5.2 设计一个去重表

以下是 MySQL/InnoDB 的教学表结构:

1
2
3
4
5
6
CREATE TABLE consumed_event (
consumer_name VARCHAR(64) NOT NULL,
event_id VARCHAR(64) NOT NULL,
consumed_at TIMESTAMP NOT NULL DEFAULT CURRENT_TIMESTAMP,
PRIMARY KEY (consumer_name, event_id)
) ENGINE=InnoDB;

联合主键中的 consumer_name 代表业务处理方,例如 inventory-service。

为什么不只用 event_id?同一订单事件可能需要库存、积分两个服务分别处理,库存已经执行不应该让积分误以为自己也执行过了。这个标识应当按业务语义设计,而不是每次重放随意换个名字。

5.3 正确的执行顺序

以一个需要记录消费身份的业务为例:

1
2
3
4
5
6
7
8
9
开启数据库事务
↓
插入 consumed_event(consumer_name, event_id)
↓
执行业务更新
↓
提交同一个数据库事务
↓
提交 Kafka Offset

应用伪代码:

1
2
3
4
5
6
7
8
9
10
11
try:
在同一个数据库事务内:
插入去重记录
执行业务更新
catch 去重表指定唯一键冲突:
回滚本次事务
确认这是已处理事件,按成功路径继续
catch 其他异常:
回滚事务并进入重试,不提交该事件之后的消费进度

业务已成功或已确认重复后,推进该分区的连续完成进度

这里必须精确识别去重表指定约束的重复,不能把所有数据库异常都当作“已处理”。业务更新影响行数、库存条件和状态转换也必须校验;更新失败应让整个事务回滚。

5.4 为什么事务必须把两者包住?

故障时机 去重记录与业务同事务时的结果
插入去重记录后,业务失败 一起回滚,下一次仍能尝试
事务成功后,Offset 提交失败 重放触发唯一约束,避免再次执行业务
两个请求同时处理同一事件 唯一约束协调竞争,成功结果只落一份

这套推导的前提是:去重与业务写入发生在同一个支持事务的数据库中,提交进度不会越过失败记录,而且去重记录保留时间覆盖重试与重放窗口。

5.5 Redis SETNX 为什么不能单独解决?

如果先执行 SETNX 记录“已处理”,再扣库存,中间崩溃后,重试可能被那个标记挡住,库存却没有扣。

反过来先扣库存再写标记,也有重复扣减的窗口。分布式锁能控制同时执行,但锁过期、释放或进程崩溃后,不能独自证明业务已经永久成功。

发短信、调用外部支付等副作用同样无法被本地数据库事务自动包住,需要下游支持幂等请求键,或额外的任务记录和对账补偿。


六、Kafka 事务:原子地提交 Kafka 内部的工作

6.1 哪类问题适合 Kafka 事务?

例如消费订单事件,计算一条风控结果,再写入另一个 Topic:

1
2
3
4
5
读取 order-events
↓
计算 risk-result
↓
写入 risk-events + 提交输入消费进度

如果输出已经写入,但输入进度没提交,重启就可能再生成一条输出。Kafka 事务可以把输出记录和输入消费组进度放在同一个 Kafka 事务中提交。

6.2 一个事务处理周期

1
2
3
4
5
6
7
8
初始化事务生产者:initTransactions()

每次处理:
poll 输入记录
beginTransaction()
send 输出记录
sendOffsetsToTransaction(本批下一位置, consumer.groupMetadata())
commitTransaction()

这是流程示意,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
2
先写数据库,再发消息 → 数据库成功后,发送可能失败
先发消息,再写数据库 → 消费者已收到事件,数据库却可能失败

即使在方法上加本地数据库的 @Transactional,也不会因此让 Kafka 的确认结果与 MySQL 的事务提交变成同一个原子操作。

7.2 本地事务只做一件能保证的事

将订单记录与待发送事件放进同一个数据库事务:

1
2
3
4
5
6
7
8
9
MySQL 本地事务
├── 写入订单
└── 写入 outbox_event
↓
事务提交
↓
后台投递器或 CDC 读取 Outbox
↓
发送 Kafka,等待成功确认

这是一种基于前面故障窗口分析得到的工程方案。它把“订单存在,就应该有对应事件”的约束先落在数据库中,再通过可恢复的异步投递完成后半段。

一个简化的待发送表:

1
2
3
4
5
6
7
8
CREATE TABLE outbox_event (
event_id VARCHAR(64) NOT NULL PRIMARY KEY,
aggregate_id VARCHAR(64) NOT NULL,
event_type VARCHAR(64) NOT NULL,
payload JSON NOT NULL,
status VARCHAR(16) NOT NULL DEFAULT 'NEW',
created_at TIMESTAMP NOT NULL DEFAULT CURRENT_TIMESTAMP
) ENGINE=InnoDB;

7.3 投递器仍可能重复发送

1
2
3
4
5
发送成功
↓
将 Outbox 标记为 SENT 之前崩溃
↓
恢复后重新发送同一事件

因此每次投递复用同一个 event_id,消费端继续做去重。Outbox 帮助恢复未完成投递,不能消除所有重复。

这个简化表结构还不是完整投递器。多实例抢占、超时重试、失败记录、清理策略都需要设计。如果同一订单要求严格顺序,投递器还要保证该订单事件的发送顺序;相同 Kafka Key 无法修正源头已经发反的顺序。


八、失败重试:不能让后续 Offset 越过错误

8.1 可恢复错误与不可恢复错误

失败 一种处理思路
下游暂时超时 有限次数重试,保持分区完成边界
JSON 格式错误 记录错误原因,转入人工检查或错误 Topic
业务数据暂未到齐 根据业务设计延后处理与补偿
认证、配置长期错误 报警并修复,避免无限空转

一个始终失败的事件可能挡住整个分区。如何绕过它,是业务语义选择:如果它是订单“创建”事件,后面的“支付”事件能否独立执行?

8.2 失败消息转移也有提交窗口

把失败事件发到 order-events-dlt 后,才能考虑推进原分区进度,而且要保留原 Topic、Partition、Offset、eventId 与错误原因。

直接发送错误消息、再提交输入进度,同样有双写窗口。在 Kafka 内部可以用事务把两者结合;采用普通发送时,需要确认投递成功并允许错误消息重复,同时做补偿与监控。

把消息移到错误 Topic 代表它被隔离了,不代表原业务已经完成。

对于严格顺序业务,将失败事件送进独立重试 Topic、让后续事件继续,可能打乱原分区的业务执行顺序,需要另外设计阻塞、版本校验或按实体恢复机制。


九、把可靠性变成可验证的实验

不用一开始就搭复杂集群,可以先验证消费者这一侧:

  1. 在测试业务完成后、提交 Offset 前人为中断进程,观察重启后的重复。
  2. 为事件添加稳定的 eventId,将去重与业务写入放入同一个数据库事务,再做同样的中断。
  3. 模拟某条记录失败,确认代码不会继续提交越过这条记录的进度。
  4. 保留一条待发送 Outbox,模拟投递失败,再观察恢复后能否继续发送。

多副本故障实验则需要多个 Broker:观察 ISR 缩小、Leader 切换,以及同步副本不足时生产者是否正确保留并恢复失败事件。

最后用这张职责表检查设计是否缺了一层:

机制 主要负责什么
副本 + ACK + 最低 ISR Kafka 存储链路的冗余与确认条件
生产者幂等 协议级发送重试的重复写入
业务事务后提交 Offset 处理失败后的重新尝试
消费业务幂等 重复到达时避免重复业务结果
Kafka 事务 Kafka 内输出记录与输入进度的原子提交
Outbox + 恢复投递 数据库提交到消息发送之间的双写窗口

学习 Kafka 的可靠性,最有价值的练习就是选一个故障时机,把数据、事件和 Offset 各自的状态写出来,再检查恢复后会不会漏处理或重复产生业务效果。