Kafka 的持久化与 ACK 机制解决了「消息不丢」,但重复是常态:生产端重试可能多写、消费端崩溃重启可能多读、事务回滚可能重放。可靠性工程的核心不是「消灭重复」,而是让重复无害化。本文系统讲解三种投递语义、重试与指数退避、幂等消费的去重手段,以及死信队列(DLQ) 这一最后防线的设计与治理。
1. 三种投递语义
1.1 语义对比
| 语义 | 保证 | 重复/丢失 | 实现成本 |
|---|---|---|---|
| At-most-once | 最多一次 | 可能丢、不重复 | 最低 |
| At-least-once | 至少一次 | 不丢、可能重复 | 中(默认) |
| Exactly-once | 精确一次 | 不丢不重 | 最高(事务/幂等) |
1.2 Kafka 各环节能保证什么
生产端:
幂等生产者 + acks=all → 不丢、单实例重试不重复(at-least-once)
存储端:
3 副本 + ISR → 已提交消息不丢(持久化)
消费端:
手动提交偏移量 → 处理成功后提交 = at-least-once
处理前提交(auto.commit) → 可能丢 = at-most-once
事务 + 幂等表 → exactly-once
1.3 选择原则
非关键/可容忍丢失:at-most-once(省成本)
默认选择:at-least-once + 幂等消费(性价比最高)
强一致关键业务:exactly-once(事务,见 kafka-transactions 篇)
一句话:Kafka 默认给你的就是 at-least-once——不丢但会重复;可靠性工程的重心 = 让重复无害(幂等)而不是幻想不重复。
2. 至少一次的重复来源
2.1 重复从哪来
① 生产端重试:网络抖动 → 同一消息 broker 收两次(幂等生产者可解)
② 消费端提交失败:处理完但 offset 没提交 → 重启重读已处理消息
③ 事务回滚重放:事务失败重来 → 消费者看到两次
④ 分区重平衡:Consumer 被踢出再入组 → 重新消费部分消息
2.2 提交策略与重复窗口
| 提交策略 | 重复窗口 | 风险 |
|---|---|---|
auto.commit=true | 无(可能丢) | 处理中崩溃丢消息 |
commitSync 处理前 | 小 | 处理失败但已提交 |
commitSync 处理成功后 | 宽 | 处理成功但未提交 → 重读 |
「处理成功后提交」是 at-least-once 的标准做法:牺牲「小概率重复」,换「不丢」。
2.3 案例分析:订单消费者
while (true) {
ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(500));
for (ConsumerRecord<String, String> r : records) {
process(r); // 处理业务(如写库)
}
consumer.commitSync(); // 全部成功后再提交偏移量
// 若 process 抛异常:不提交 → 重启后重读该批 → 重复处理
}
一句话:at-least-once 的重复窗口由**「处理成功 → 提交偏移量」**之间的崩溃决定——这个窗口就是你要用幂等兜住的范围。
3. 幂等消费:让重复无害化
3.1 幂等 = 同一输入重复处理结果一致
业务操作天然幂等的例子:UPDATE SET status='paid' WHERE id=123(结果相同);天然不幂等的例子:INSERT 新记录、余额累加、发短信。
3.2 幂等表(Idempotency Table)
对非幂等操作,用「唯一键 + 去重表」把重复变无害:
-- 消息去重表
CREATE TABLE processed_events (
msg_id VARCHAR(64) PRIMARY KEY, -- 消息唯一 ID(业务幂等键)
topic VARCHAR(64) NOT NULL,
partition INT NOT NULL,
offset BIGINT NOT NULL,
processed_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP,
UNIQUE KEY uk_po (topic, partition, offset) -- 可选:按位置去重
);
// 消费端:先插幂等表再处理业务
public void process(ConsumerRecord<String,String> r) {
try {
dao.insertProcessed(r.key()); // ① 插入幂等表(唯一键)
applyBusinessLogic(r); // ② 业务逻辑(同一事务)
consumer.commitSync(); // ③ 提交
} catch (DuplicateKeyException e) {
consumer.commitSync(); // 已处理过 → 跳过并提交
}
}
关键:insertProcessed 与 applyBusinessLogic 要在同一数据库事务里,否则去重标记与业务结果可能不一致。
3.3 业务键 vs 消息位置
| 去重依据 | 适用 | 例子 |
|---|---|---|
| 业务幂等键 | 同一业务事件多次投递 | 订单号、交易号 |
| 消息位置 (topic,partition,offset) | 消费端重读 | 配合偏移量 |
推荐双管齐下:业务键防「生产端多写」,消息位置防「消费端重读」。
3.4 用 Redis 做去重(高吞吐轻量)
// 用 SETNX 加唯一键,TTL 覆盖重试窗口
boolean first = jedis.setnx("event:" + msgId, "1");
jedis.expire("event:" + msgId, 3600); // 1 小时窗口
if (!first) { /* 已处理,跳过 */ }
Redis 去重快但不持久(崩溃丢标记),适合配合 DB 幂等表做第一道过滤。
一句话:幂等消费 = 用「业务唯一键」或「消息位置」去重——DB 幂等表可靠、Redis 快、两者配合最佳;核心是「去重标记与业务结果同事务」。
4. 重试与指数退避
4.1 为什么需要重试层
消费失败的原因常是瞬时的:下游数据库抖动、依赖服务 503、网络闪断。直接丢给 DLQ 太粗暴,先重试更合理。
4.2 重试层次
① 客户端自动重试:broker 侧超时/可重试错误(生产端 retries)
② 消费端代码重试:单消息处理失败 → 重新 poll 重试
③ 重试 Topic 模式:失败消息投递到「重试队列」,延迟后回读
4.3 指数退避(Exponential Backoff)
int attempts = 0;
long baseDelay = 1000; // 1s 起步
while (attempts < 5) {
try {
process(msg);
break;
} catch (RetryableException e) {
attempts++;
long delay = baseDelay * (1L << (attempts - 1)); // 1,2,4,8,16s
delay += new Random().nextInt(500); // 加抖动防惊群
Thread.sleep(delay);
}
}
// 5 次仍失败 → 投递 DLQ
加抖动(jitter) 很重要:多个消费者同时重试,不加抖动会形成「重试风暴」。
4.4 重试 Topic 模式
高级做法:失败消息写入 orders.retry(带延迟消费),由专门消费者延迟 N 秒后回读原 Topic 或直接处理:
orders 主 Topic → 失败 → orders.retry(延迟 30s 后回读)
重试 3 次仍失败 → orders.dlq(死信)
一句话:重试 = 指数退避 + 抖动 + 次数上限——先给瞬时故障「自救机会」,救不动再进 DLQ;别无限重试拖垮整个消费链。
5. 死信队列(DLQ)设计与治理
5.1 DLQ 是最后防线
DLQ 收容重试后仍失败的坏消息,避免「一条毒消息卡死整个分区消费」。分区内消息是顺序处理的,若失败消息不摘出来,后续消息全部阻塞——DLQ 就是「拔掉毒刺」。
5.2 DLQ 设计要点
命名规范:<topic>.dlq 或 <topic>.<group>.dlq
消息内容:保留原始 payload + header(原因、重试次数、原始 offset)
保留策略:长保留(审计用),别短于业务回溯周期
监控:DLQ 深度/Lag 是核心告警
// 消费端:重试耗尽后投递 DLQ
catch (Exception e) {
if (attempts >= MAX_ATTEMPTS) {
Map<String, String> headers = Map.of(
"err-reason", e.getMessage(),
"err-retries", String.valueOf(attempts),
"err-original", r.topic() + ":" + r.partition() + ":" + r.offset()
);
dlqProducer.send(new ProducerRecord<>("orders.dlq",
r.key(), r.value(), headers));
consumer.commitSync(); // 摘除毒消息,继续消费
}
}
5.3 DLQ 治理闭环
① 落 DLQ(带原因/重试次数)
② 告警(DLQ 深度 > 阈值)
③ 人工/工具排查原因(反序列化错?下游 bug?)
④ 修复后重新投递回主 Topic(重放)
⑤ 定期清理过期 DLQ 消息
5.4 DLQ 常见误用
| 误用 | 问题 |
|---|---|
| 无重试直接进 DLQ | 瞬时故障也进,DLQ 爆炸 |
| DLQ 无限保留 | 存储成本失控 |
| 无告警 | 毒消息静默堆积 |
| 重放不做幂等 | 重放又造重复 |
一句话:DLQ = 「毒消息隔离区」——让一条坏消息不阻塞整个分区;设计上保留原因与原始位置、设深度告警、走重放闭环,DLQ 才能真正成为兜底而非藏污。
6. 端到端可靠性架构
6.1 全链路可靠组合
生产端:幂等生产者 + acks=all + min.insync=2 → 不丢不重(单实例)
存储端:3 副本 + ISR → 持久可靠
消费端:手动提交 + 幂等表 + 重试退避 + DLQ → 重复无害 + 毒消息隔离
6.2 可靠性设计决策表
| 组件 | 默认 | 强化 |
|---|---|---|
| acks | 1 | all(金融级) |
| min.insync.replicas | 1 | 2 |
| 消费提交 | auto | 手动 + 成功后才提交 |
| 去重 | 无 | 业务键幂等表 |
| 重试 | 无 | 指数退避 + 抖动 |
| 失败处理 | 丢弃 | DLQ + 告警 + 重放 |
6.3 监控可靠性指标
生产端:发送失败率、重试率、delivery.timeout
消费端:消费 Lag、处理失败率、重试次数分布
DLQ:深度、增长速率、重放成功率
一句话:端到端可靠 = 幂等生产 + 持久存储 + 幂等消费 + 重试 + DLQ 的组合拳——每一层解决一个「不可靠」,串起来才是可靠。
7. 常见坑与最佳实践
7.1 常见坑
| 坑 | 现象 | 对策 |
|---|---|---|
| 无幂等靠「应该不会重复」 | 线上偶发重复事故 | 业务键幂等表 |
| 失败消息不摘 DLQ | 分区消费卡死 | 重试耗尽进 DLQ |
| 无限重试 | 消费链雪崩 | 次数上限 + 退避 |
| 去重标记与业务不同事务 | 去重失效 | 同事务落幂等表 |
| DLQ 无告警 | 静默堆积 | 深度/增长告警 |
| 重放不幂等 | 重放再造重复 | 重放走幂等键 |
7.2 最佳实践清单
- 默认 at-least-once + 手动提交,业务键幂等兜底;
- 重试用指数退避 + 抖动 + 上限,别无限;
- DLQ 带原因/重试次数/原始位置,深度告警;
- 去重与业务同事务,别半吊子;
- 定期演练重放,验证 DLQ 重放闭环;
- 监控四指标:失败率、Lag、重试分布、DLQ 深度。
8. 总结
本文把 Kafka 可靠性工程梳理成完整体系:
| 环节 | 核心 |
|---|---|
| 语义 | at-least-once 默认,重复是常态 |
| 重复来源 | 生产重试 / 消费重读 / 重平衡 |
| 幂等 | 业务键 / 消息位置,去重与业务同事务 |
| 重试 | 指数退避 + 抖动 + 次数上限 |
| DLQ | 毒消息隔离 + 原因留痕 + 重放闭环 |
| 架构 | 幂等生产 + 幂等消费 + DLQ 组合拳 |
一句话记住:Kafka 可靠性不是「消灭重复」,而是**「让重复无害 + 让失败可治理」**——幂等表去重、退避重试、DLQ 兜底三层防线,配告警与重放闭环。默认 at-least-once,关键业务上事务,失败别硬扛进 DLQ——可靠是设计出来的,不是祈祷出来的。
继续阅读
探索更多技术文章
浏览归档,发现更多关于系统设计、工具链和工程实践的内容。