极高 进阶
消息丢失与可靠性保证#
一句话答案#
三环节保证不丢:生产端(acks=all + 重试)→ Broker(多副本 ISR + min.insync.replicas)→ 消费端(手动提交 offset + 幂等处理)。任何一环断裂都会丢消息。
核心要点
全链路分析#
Producer ──发送──→ Broker ──存储──→ Consumer
│ │ │
├── 网络丢失 ├── 宕机丢失 ├── 自动提交后处理失败
├── 超时未重试 ├── 副本不足 └── offset 跳过
└── 序列化失败 └── 刷盘前崩溃plaintext第一环:Producer → Broker#
acks 参数详解:
| acks | 含义 | 性能 | 可靠性 |
|---|---|---|---|
| 0 | 不等待任何确认 | 最高 | 最低(发完就忘) |
| 1 | 等待 Leader 确认 | 中等 | 中等(Leader 宕机可能丢) |
| all(-1) | 等待 ISR 中所有副本确认 | 最低 | 最高 |
acks=all 的陷阱:
如果 ISR 中只有 Leader 一个副本(其他副本被踢出 ISR)
→ acks=all 效果等于 acks=1
→ Leader 宕机仍会丢消息
解决:配合 min.insync.replicas=2
→ ISR 副本数 < 2 时,acks=all 的写入直接报错(NotEnoughReplicas),acks=0/1 不受此参数约束
→ 用可用性换可靠性plaintext推荐配置:
# Producer 端(Kafka 3.0 起 acks=all、enable.idempotence=true 已是默认值)
acks = all
enable.idempotence = true # 幂等生产者:重试不重复,且 in.flight ≤ 5 时仍保证单分区有序
retries = 2147483647 # 默认值即 Integer.MAX_VALUE,实际重试时长由 delivery.timeout.ms(默认 120s)限制
retry.backoff.ms = 100 # 重试间隔
max.in.flight.requests.per.connection = 5 # 开幂等后不必再设 1;没开幂等时要设 1 才能防重试乱序
# Broker 端
min.insync.replicas = 2 # ISR 至少 2 个副本
replication.factor = 3 # 每个 Partition 3 个副本properties第二环:Broker 存储#
Kafka 消息写入流程:
Producer → Leader Broker (写入 Page Cache) → 操作系统异步刷盘
↓
Follower 从 Leader 拉取复制plaintext可能丢失的点:
- Page Cache 未刷盘时 Broker 宕机 → 配合副本机制,其他副本有数据
- 所有 ISR 副本同时宕机 → 极端情况,概率极低
不建议用 flush.messages=1(每条消息同步刷盘):
- 性能下降 10 倍以上
- Kafka 的设计哲学是用副本保可靠性,而不是用刷盘
第三环:Consumer 消费#
自动提交 offset 的坑:
Consumer 拉取消息 → 处理失败但异常被吞掉 / 交给异步线程处理
→ 下次 poll() 时自动提交了这批 offset → 消息不会再投递 → 消息"丢失"(实际是跳过了)
(同步处理完再 poll,自动提交也是 at-least-once;丢消息出在「提交早于处理完成」)plaintext正确做法:手动提交
ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(100));
for (ConsumerRecord<String, String> record : records) {
// 1. 先处理业务逻辑
processMessage(record);
// 2. 处理成功后再提交这一条的 offset(提交的是「下一条要读的位置」= offset + 1)
consumer.commitSync(Map.of(
new TopicPartition(record.topic(), record.partition()),
new OffsetAndMetadata(record.offset() + 1)));
}
// 注意:无参 commitSync() 提交的是本次 poll 返回的全部位点,
// 放在循环里处理完第 1 条就会把整批都标成已消费,后面几条处理失败就丢了。
// 按批提交时应在整个 for 循环结束后再调用无参 commitSync()。java手动提交的粒度选择:
| 方式 | 特点 |
|---|---|
| 每条提交 | 最安全,但性能差(频繁 IO) |
| 每批提交 | 性能好,但失败时整批重试(需要幂等) |
| 异步提交 | 性能最好,但提交失败不会阻塞(需要回调处理) |
推荐: 每批处理完同步提交 + 消费端做幂等处理。
消费端幂等#
手动提交保证不跳过消息,但可能重复消费(处理完了但提交 offset 前宕机)。
幂等方案:
- 数据库唯一键(相同消息 ID 插入报唯一冲突,忽略)
- Redis SET NX(消息 ID 做 Key,消费前检查是否存在)
- 业务状态机(订单已支付 → 重复支付消息忽略)
完整的可靠性配置清单#
Producer:
✅ acks = all
✅ retries 保持默认(Integer.MAX_VALUE,受 delivery.timeout.ms 限制)
✅ enable.idempotence = true(防止重试导致的重复消息;3.0 起默认开启)
Broker:
✅ replication.factor = 3
✅ min.insync.replicas = 2
✅ unclean.leader.election.enable = false(不允许非 ISR 副本成为 Leader)
Consumer:
✅ enable.auto.commit = false
✅ 手动 commitSync / commitAsync + callback
✅ 消费逻辑做幂等plaintext面试回答(2分钟版)
消息不丢需要三环节同时保证。Producer 端:acks=all 确保 ISR 所有副本确认,配合 retries 自动重试网络瞬时故障。但 acks=all 有个陷阱——如果 ISR 只有 Leader 一个副本,效果等于 acks=1,所以 Broker 端必须配 min.insync.replicas=2,ISR 不足时直接拒写。Broker 端还要配 replication.factor=3、unclean.leader.election=false 防止非 ISR 副本当 Leader 导致数据丢失。Consumer 端关掉自动提交,处理完业务逻辑再手动 commitSync。但手动提交可能导致重复消费(处理完了没来得及 commit 就宕机),所以消费端必须做幂等——用数据库唯一键或 Redis SETNX。三环节加起来是 At-Least-Once + 业务幂等,这是大多数生产系统的选择。
追问与易错
追问方向:
- “acks=all 影响性能多少?”→ 延迟增加(等待副本确认),吞吐量下降 20%-40%,但可通过 batch 和异步发送缓解
- “消费者宕机未 ACK 怎么办?”→ 消息会重新投递给同组其他消费者,需要幂等
- “min.insync.replicas 设太大会怎样?”→ 可用性降低,ISR 副本不够时 acks=all 的写入全部拒绝(读不受影响)
易错点:
- ❌ “acks=all 就绝对不丢”——ISR 只有 Leader 时等于 acks=1,必须配 min.insync.replicas
- ❌ “消费完就 commit offset”——应该处理完业务逻辑再 commit
- ❌ “Kafka 不丢消息”——任何中间件都可能丢,关键是配置正确 + 业务兜底