Kafka 如何保证消息不丢失
从生产者、Broker 和消费者三段链路分析 Kafka 消息丢失的原因,以及 acks、ISR、副本、手动提交 Offset 和业务幂等的完整保障方案。
Kafka 如何保证消息不丢失
面试回答
Kafka 保证消息不丢失需要覆盖生产者、Broker 和消费者三段链路。生产者侧使用
acks=all、开启重试和幂等生产,并检查发送回调,对最终失败进行落库、告警或补偿;Broker 侧通常设置三副本、min.insync.replicas=2,并禁止非同步副本参与 Leader 选举,保证消息在足够多的同步副本写入后才确认成功;消费者侧关闭自动提交,在业务处理成功后再提交 Offset,失败时不提交以便重新消费。这样通常实现至少一次投递,所以消费端还必须通过事件 ID、唯一索引或状态机保证幂等。涉及数据库与 Kafka 的一致性时,可以使用本地消息表或 Transactional Outbox。Kafka 的消息不丢失是端到端保障,不是只配置一个参数。
一句话总结:
生产者确认成功,Broker 多副本可靠保存,消费者处理成功后提交,业务端用幂等和补偿兜底。
详细讲解
Kafka 消息不丢失不是靠某一个参数实现的,而是一个端到端问题:
生产者必须确认消息发送成功,Broker 必须把消息可靠地保存在足够多的副本中,消费者必须在业务处理成功后再提交 Offset。
其中任何一个环节处理不当,都可能出现“业务上看起来消息丢了”的情况。
一、先明确什么叫“消息丢失”
常见的丢失场景主要有四类:
- 生产者调用发送方法后没有检查结果,实际发送失败却被业务忽略。
- Leader 收到消息后就返回成功,消息还未复制到 Follower,Leader 随后故障。
- 消费者先提交 Offset,再执行业务;业务处理失败后,消息不会再次投递。
- Kafka 中的消息没有丢,但写数据库、调用下游等业务操作失败,又没有补偿机制。
因此,Kafka 中仍能查到消息,不代表业务一定处理成功;反过来,业务没有结果,也不一定是 Broker 丢了消息。
二、生产者如何避免丢消息
1. 使用 acks=all
生产者的 acks 决定 Broker 在什么条件下确认发送成功:
| 配置 | 确认时机 | 风险 |
|---|---|---|
acks=0 | 不等待 Broker 响应 | 发送失败也无法感知,风险最高 |
acks=1 | Leader 写入本地日志后返回 | Leader 故障且副本尚未同步时可能丢失 |
acks=all | 当前 ISR 中的副本都确认后返回 | Kafka 能提供的最强生产者确认保证 |
关键配置:
acks=all
acks=all 并不等于绝对不丢。它还需要与副本数和 min.insync.replicas 配合,否则 ISR 中只剩 Leader 一个副本时,依然可能成功写入。
2. 开启重试和幂等生产
网络抖动、Leader 切换等临时故障可能导致发送失败,生产者需要允许重试:
enable.idempotence=true
acks=all
retries=2147483647
max.in.flight.requests.per.connection=5
delivery.timeout.ms=120000
幂等生产者会给消息附加 Producer ID 和序列号,使 Broker 能识别同一生产会话内的重复写入,避免因重试产生重复消息。
需要注意:
- 幂等生产解决的是重试导致的重复写入,不是发送失败后的业务补偿。
delivery.timeout.ms到期后,生产者仍可能最终失败。- 业务不能无限依赖客户端重试,最终失败必须记录、告警或进入补偿流程。
3. 必须检查发送结果
下面这种“只发送、不处理结果”的写法存在风险:
producer.send(record);
应该检查回调或 Future 的执行结果:
producer.send(record, (metadata, exception) -> {
if (exception != null) {
// 记录原始消息、告警并进入补偿流程
handleSendFailure(record, exception);
}
});
序列化失败、鉴权失败、消息过大和超时等错误,最终都需要业务明确处理。不能把“调用过 send”当成“消息已经可靠进入 Kafka”。
4. 数据库与 Kafka 的一致性
如果业务流程是:
更新数据库
→
发送 Kafka 消息
数据库更新成功、消息发送失败时,仍然会造成业务事件丢失。
常见解决方式是本地消息表,也叫 Transactional Outbox:
同一个数据库事务
├── 更新业务数据
└── 写入待发送事件表
事务提交后
→ 后台任务投递 Kafka
→ 成功后标记事件已发送
这样即使 Kafka 暂时不可用,消息也可以从本地消息表继续重试。Kafka 事务可以保证 Kafka 内部多条记录及消费 Offset 的原子性,但不会自动把外部数据库事务包含进来。
三、Broker 如何保证消息可靠保存
1. 合理设置副本数
生产环境通常使用:
replication.factor = 3
一个分区包含一个 Leader 和多个 Follower。生产者与 Leader 交互,Follower 从 Leader 复制日志。某个 Broker 故障后,可以从仍然同步的副本中选举新 Leader。
副本数提高的是容错能力,但副本只有真正保持同步才有意义。
2. 理解 ISR
ISR 是 In-Sync Replicas,即当前与 Leader 保持同步的副本集合。
例如一个三副本分区:
Leader A
Follower B
Follower C
ISR = [A, B, C]
如果 C 长时间跟不上 Leader,它会被移出 ISR:
ISR = [A, B]
acks=all 等待的是当前 ISR 的确认,而不是永远等待所有配置副本。
3. 配置 min.insync.replicas
推荐组合:
replication.factor = 3
min.insync.replicas = 2
acks = all
它表达的含义是:
至少要有两个同步副本可用,写入才允许成功。
如果 ISR 只剩一个副本,Broker 会拒绝写入。此时系统牺牲部分可用性,避免在单副本状态下继续写入并承担更高的数据丢失风险。
4. 禁止非同步副本强行成为 Leader
unclean.leader.election.enable=false
如果允许非 ISR 副本成为 Leader,在所有同步副本不可用时,Kafka 可以更快恢复分区服务,但这个副本可能缺少最新消息,从而造成数据丢失。
这是一个典型的取舍:
允许非同步副本选举:可用性更高,可能丢数据
禁止非同步副本选举:一致性更强,分区可能暂时不可用
5. acks=all 不等于每条消息都立即 fsync
Kafka 的持久性主要依赖顺序写日志、操作系统页缓存和多副本机制。Broker 返回成功,并不表示每个副本都对该消息单独执行了一次物理磁盘 fsync。
实际生产中,通常通过跨 Broker 副本降低单机和单盘故障风险,而不是强制每条消息同步刷盘。若多个副本同时发生不可恢复故障,仍不存在数学意义上的绝对零丢失。
四、消费者如何避免“消费丢失”
1. 关闭自动提交 Offset
enable.auto.commit=false
自动提交可能出现以下顺序:
拉取消息
→ 自动提交 Offset
→ 执行业务
→ 进程崩溃
重启后,消费者会从已提交的下一个 Offset 继续消费,刚才尚未处理成功的消息就被跳过了。
2. 业务成功后再提交 Offset
更可靠的顺序是:
拉取消息
→ 执行业务
→ 业务成功
→ 提交 Offset
如果业务处理失败,就不提交 Offset,让消息后续重新消费:
while (true) {
ConsumerRecords<String, String> records = consumer.poll(Duration.ofSeconds(1));
for (ConsumerRecord<String, String> record : records) {
process(record);
}
consumer.commitSync();
}
示例表达的是基本原则。实际批量消费时,如果一批消息中只有部分成功,需要准确管理各分区可提交的 Offset,不能直接把失败消息之后的位置一起提交。
3. 为什么还需要业务幂等
“处理成功后提交”可以避免消息被跳过,但会引入重复:
业务处理成功
→ 提交 Offset 前进程崩溃
→ 重启后再次消费同一条消息
因此,可靠消费通常采用至少一次投递,并在业务层保证幂等。常见方式包括:
- 使用事件 ID 建立数据库唯一索引。
- 建立消费去重表或 Inbox 表。
- 使用业务状态机,只允许合法状态转换。
- 更新时带版本号或业务条件。
- 对扣款、发券等操作使用唯一业务流水号。
不丢消息通常意味着允许重复,再通过幂等消除重复影响。
4. Kafka 内部链路的 Exactly Once
如果消费 Kafka 后仍然写回 Kafka,可以使用 Kafka 事务,把输出记录和消费 Offset 放进同一个事务:
消费消息
→ 处理
→ 写入下游 Topic
→ sendOffsetsToTransaction
→ 提交事务
下游消费者设置:
isolation.level=read_committed
这样只读取已提交事务的数据。
但如果最终写入的是 MySQL 等外部系统,Kafka 事务并不能自动保证 Kafka 与数据库的原子性,仍需要数据库事务、幂等、Outbox 或 Inbox 等方案。
五、三种投递语义
| 语义 | Offset 与业务处理顺序 | 结果 |
|---|---|---|
| At Most Once | 先提交,再处理 | 可能丢失,通常不重复 |
| At Least Once | 先处理,再提交 | 不轻易丢失,可能重复 |
| Exactly Once | 事务或幂等机制协调 | 业务效果恰好一次 |
绝大多数业务系统采用:
Kafka 至少一次投递 + 消费端业务幂等。
这通常比追求所有环节的强事务更容易实现,也更便于扩展。
六、常见故障与保障措施
| 故障场景 | 可能后果 | 主要保障 |
|---|---|---|
| 生产者网络抖动 | 发送失败或结果未知 | 重试、幂等生产、回调处理 |
| Leader 写入后立即故障 | 未同步消息丢失 | acks=all、副本、ISR |
| ISR 只剩一个副本 | 单点故障风险升高 | min.insync.replicas 拒绝写入 |
| 同步副本全部不可用 | 错误选主后丢数据 | 禁止 unclean leader election |
| 消费者先提交后处理 | 业务消息被跳过 | 关闭自动提交,成功后提交 |
| 处理成功但提交前崩溃 | 重复消费 | 业务幂等、唯一流水号 |
| 数据库成功、Kafka 失败 | 业务事件未发出 | Transactional Outbox |
| Kafka 成功、下游失败 | 业务结果缺失 | 重试、死信、补偿、告警 |
七、常见误区
误区一:配置 acks=all 就绝对不会丢
acks=all 只解决生产者到 Broker 的确认强度,还要配合副本数、min.insync.replicas、正确选主和消费者提交策略。
误区二:配置无限重试就不会丢
重试仍然受 delivery.timeout.ms 等条件限制,而且权限错误、序列化错误等问题无法靠盲目重试解决。最终失败必须被业务感知和补偿。
误区三:手动提交 Offset 就是 Exactly Once
手动提交只能帮助建立“处理成功后提交”的顺序,崩溃窗口仍可能造成重复消费,因此还需要业务幂等。
误区四:Kafka 事务可以自动覆盖数据库
Kafka 事务主要协调 Kafka 内部的消息与 Offset,不会自动与 MySQL 等外部数据库形成同一个原子事务。
误区五:Broker 返回成功就等于所有副本已经物理刷盘
Kafka 的确认、副本复制和磁盘刷盘是不同概念。高可靠主要依靠多副本和同步副本约束,而不是把每条消息都单独同步刷盘。
八、监控和运维同样重要
配置正确后,还应持续监控:
- 发送失败率、重试次数和发送延迟。
- ISR 收缩、未充分复制分区和离线分区。
- Broker 磁盘空间、磁盘延迟和网络异常。
- 消费积压、消费失败、重平衡次数。
- 死信消息和补偿任务堆积。
- 业务事件与最终结果的对账差异。
没有监控和对账,即使消息已经丢失,系统也可能长期无法发现。
评论与讨论
回复 :
留下你的想法