“消息处理失败了,扔进 Dead Letter Queue,等 bug 修完 replay 回去。这大概是消息系统里听起来最简单的一个环节。 大多数教程讲到配一个 DLQ topic就收工了。但我维护过的几个消息系统里,跟 DLQ 相关的 on-call incident,比 consumer lag 导致的还多。问题不在消息怎么进 DLQ, 问题在消息怎么从 DLQ 出来。”
01 Poison pill 还是 transient failure
区分 poison pill 和 transient failure 不是一个简单的 if/else。在 Kafka batched consumer 里,一条消息的异常类型、batch commit 语义、和 partition coupled together,让这个判断变成一个三维取舍。
Kafka consumer poll() 一次拿 max.poll.records 条消息(默认 500)。你拿到 500 条,逐条处理,第 37 条 throw exception。这条消息可能是三种 failure 类型之一:
Deserialization failure:Spring Kafka 的 ErrorHandlingDeserializer 能捕获这种异常(注意这是 Spring Kafka 的类,不是 Apache Kafka 自带的)。消息连deserialization都过不了,绝大多数情况是 poison pill,直接 DLQ。但这条规则有一个必须排掉的反例:schema registry 不可达时,Confluent 的 Avro / Protobuf deserializer 一样抛 SerializationException,这是纯 transient 的,而且它会让整个 topic 的所有消息同时“deserialization失败”,全量灌进 DLQ。所以”deserialization失败 = poison pill”要加一个前置判断:registry 通不通。
Downstream timeout / 503:下游临时挂,retry 有意义。但 retry 期间这个 partition 的消费被阻塞(除非你把失败消息发到 retry topic,那又跳到 retry ordering 问题)。
business logic exception on valid schema:比如
amount: -1或currency: null。deserialization成功了(schema 合法),但业务逻辑永远会 throw NPE / IllegalArgument。你的catch (Exception e)看到的是一个 NPE,跟下游 timeout 的 exception 长得不一样,但你的 retry 逻辑如果只数 retry count 不看 exception type,它会重试 10 次然后扔 DLQ,用了 10 次无意义的重试才做到跟直接 DLQ 一样的结果,期间阻塞了 partition。
真实做法是按 exception type 做分类树:DeserializationException / IllegalArgumentException / NullPointerException(语义层 poison pill)> 立刻 DLQ;RetriableException / TimeoutException / ConnectException > retry with backoff。但这个分类树的维护本身是个持续工程,每加一个新的下游依赖,分类树就得更新。
02 Retry topology
retry topology 有三种主流选择,每种在 Kafka 的 partition ordering 语义下都会丢掉一样东西,ordering、throughput、或 auto-recovery 能力。没有”正确选择”,只有”你选了丢哪一头”。
方案 A:consumer 内 blocking retry。Partition ordering 完好,因为 partition 被阻塞了,后续消息等着。代价:一条 poison pill 卡住整个 partition 几十分钟。如果 blocking 时间超过 max.poll.interval.ms(默认 300 秒),consumer 自己会发现两次 poll 之间超时了,主动发 LeaveGroup 退出 group,注意不是 broker 判定它死了:KIP-62 之后 heartbeat 跑在独立的后台线程,你 block 在处理逻辑里 broker 照样收得到心跳,session.timeout.ms 根本不会触发。退组触发 rebalance,partition 重新分配,这一批未提交的 offset 全部作废 > 重复消费。
方案 B:tiered retry topics(orders.retry-1 delay 1min > orders.retry-2 delay 5min > orders.retry-3 delay 30min > DLQ)。Spring Kafka 的
@RetryableTopic一个 annotation 就配好。不阻塞原 partition。但:消息 N 被发到 retry topic 后,N+1 到 N+100 继续消费。N 从 retry topic 回来时,它在时间线上的位置已经变了。如果业务依赖 per-key ordering(比如同一个账户的余额 +100 > -50 > +30 必须按顺序),这条路就碎了。额外代价:每个 source topic 多 3-4 个 retry topic + 各自的 listener container,monitoring 和 alerting 的配置量翻几倍。而且这些 container 默认跟主 listener 在同一个 consumer group 里,不是各自独立的 group,所以任何一个 retry topic 上的 rebalance 都会连带触发主 topic 的无谓 rebalance,一个 annotation 悄悄给你的主链路加了几个抖动源。方案 C:直接 DLQ,不重试。最简单,partition ordering 不受影响(失败消息被拿走,后续继续)。但所有 transient failure 丧失自愈能力,一个下游重启 30 秒导致的 timeout 消息,本来 retry 一次就好,现在要等人手动 replay。
Kafka vs SQS vs RabbitMQ 简单说就是:SQS 有 maxReceiveCount + DLQ redrive 是 native feature;RabbitMQ 有 TTL-based dead-letter exchange + x-death header 自动计重试次数(但”计到第几次转终极 DLQ”这个判断还是要你自己写)。Kafka 这层全靠你自己建或者用 Spring Kafka 的 wrapper,但 wrapper 的默认行为(比如 @RetryableTopic 的 autoCreateTopics、以及控制”相同 backoff 间隔要不要复用同一个 topic”的 sameIntervalTopicReuseStrategy,3.x 之前叫 fixedDelayTopicStrategy)要读源码才知道语义。
03 Bug 修完了,replay DLQ:三个你没想过的问题
bug 修完了 replay DLQ 听起来就是”把消息 produce 回原 topic”,但 replay 时你会撞上三个设计 consumer 时根本没考虑过的问题:ordering violation、partial side effect 重放、和从未规划过的 consumer idempotency requirement。
Ordering violation:DLQ 里的消息来自多个 partition,跨好几天。replay 时用原 key produce 回原 topic,消息落到正确的 partition(前提是 partition 数没扩过;扩过容的 topic,key > partition 的映射已经变了,replay 会落到另一个 partition,ordering 碎得更彻底),但 offset 在最新位置,跟当天的 live traffic 混在一起。如果 consumer 有 per-key 的时序假设(比如按 event_time 做聚合、或依赖 offset 递增等价于时间递增),3 天前的消息出现在今天的 offset 位置,聚合结果会出错。
Partial side effects:消息第一次处理时可能已经完成了部分 side effect,发了 notification email、写了 audit log、扣了库存,只是在最后一步(比如更新 order status)时 fail 了才进的 DLQ。replay 重跑整个 handler:email 发两遍、库存扣两次。你的 handler 大概率没有 per-step 的 checkpoint。
Consumer idempotency 的”发现时刻”:
enable.auto.commit=false+ 手动commitSync给你的是 at-least-once,不是 exactly-once。处理完了、commit 之前 crash,接管的 consumer 从上次提交的位点重来,这一批全部重放。也就是说重复消费在正常流程里就已经会发生了,rebalance 一次、crash 一次就来一遍。但大多数 consumer 根本没做 per-message idempotency check,不是因为不需要,而是因为平时 rebalance 不频繁,重复量小到没人注意,对账差异被当成噪声抹平了。DLQ replay 只是把这件事放大到没法忽略:一次 replay 成千上万条消息第二次进 consumer,之前被噪声掩盖的不幂等,这时候变成账面上对不上的数字。大多数 team 发现自己需要 consumer-level idempotency,是在第一次 replay DLQ 的时候,但这个需求一直都在。真想要 exactly-once,得上 Kafka 事务(transactional.id + sendOffsetsToTransaction + 消费端 isolation.level=read_committed),或者干脆把 consumer 做成幂等的。
04 Silent-DLQ failure:最危险的 DLQ 故障是没人看它
DLQ 最常见也最致命的 failure mode 不是 replay 出错或消息太多,是根本没人看它。消息在 DLQ 里安静地腐烂几周,直到下游报表对不上才被发现。
场景:一个支付确认 consumer 每天约 0.1% 的消息 fail(downstream timeout、偶尔的 schema mismatch),默默进 DLQ。每天约 200 条。没人设 DLQ depth alert。三周后有人查报表发现差了几千条,DLQ 里堆了 4200+ 条消息,大部分是 transient failure 早就可以 retry 了的,其中一部分是 time-sensitive 的 payment confirmation。
为什么没人看:SQS 有 ApproximateNumberOfMessagesVisible metric 可以直接配 CloudWatch alarm。但 Kafka 的 DLQ 就是一个普通 topic,如果没有 consumer 在消费它,就没有 consumer lag metric。它不会自己报警。你需要单独跑一个 monitor consumer、或用 Kafka Exporter / Burrow 主动监控这个 topic 的 high watermark。
必须设的三个告警:DLQ depth > 0;oldest message age > 1 hour;growth rate > 50/min。
DLQ 应该是 triage queue(进来就分拣、处理、清空),不是 permanent parking lot。
05 Partial-batch failure:message 37 of 500 挂了,你的 offset 怎么 commit
Kafka 的 offset commit 是 per-partition 的一个水位线,它没有”ack 这 500 条里的 498 条、skip 其中 2 条”的原语。(水位线本身不要求单调递增,seek 回去 commit 一个更小的 offset 完全合法,kafka-consumer-groups --reset-offsets 就是这么干的;但它始终只是一个位点,不是逐条 ack。)partial-batch failure 时你的每个选择都在丢一样东西。
max.poll.records=500,poll() 拿到 offset 0-499 这 500 条。逐条处理,offset 37 这条 fail。
选项 A:
consumer.commitSync(Collections.singletonMap(partition, new OffsetAndMetadata(37))),这里有两个必须踩准的细节。第一,Kafka 的 committed offset 语义是”下一条要读的消息的 offset“,要从 37 重来就提交 37,不是 36;提交 36 会把已经成功的 36 再重放一遍。第二,光 commit 不够:commit 不改变 consumer 的内存 position。poll() 返回 500 条之后 position 已经在 batch 末尾了,committed offset 只决定 rebalance 或重启后从哪儿恢复;同一个 consumer 继续跑,下一次 poll 照样从 500 往后拿。要当场重来必须再调consumer.seek(partition, 37)。停止当前 batch 后,38-499 这 462 条白处理了。如果你的 consumer 不 idempotent,重新消费 38-499 会重复执行 side effect。如果每条处理 10ms,浪费约 4.6 秒重处理时间。选项 B:catch 37 的 exception,continue loop,处理 38-499,commit offset 500(同样是”下一条要读的”语义,500 表示 0-499 全部处理完)。消息 37 永远丢失(offset 已经推过它了)。除非你在 catch 里手动
producer.send(dlqTopic, failedMessage)把它发到 DLQ,但这个 produce 本身也可能 fail。选项 C:维护一个 per-partition 的 pending-failure list,batch 结束后先把 failed messages produce 到 DLQ,确认 produce 成功后再 commit offset 500。最安全,但引入了 producer + consumer 之间的 transactional 耦合,如果 DLQ produce 成功但 offset commit fail(consumer crash at this exact moment),下次 poll 会从上一次成功提交的位点重来,也就是整个 batch 从头再消费一遍,其中 37 已经在 DLQ 里了(DLQ 里重复一条)。Kafka 事务给了这一步一个现成原语:把 DLQ 的 produce 和
sendOffsetsToTransaction塞进同一个事务,要么都成立要么都不成立,代价是消费端要开 read_committed,吞吐和延迟都得交学费。
enable.auto.commit=true 的陷阱:很多人以为 auto-commit 是个后台定时器,其实它是在你调用 poll() 的时候执行的,poll() 一进去先检查距上次提交是否超过 auto.commit.interval.ms(默认 5000ms),超过就把上一次 poll 返回的最大 offset 提交掉。所以坑不在”你正在处理第 37 条时被偷偷提交了”(你在处理循环里的时候没人在提交),而在于:你 catch 住 37 的异常、继续跑完这一批、下一次 poll() 一进去,offset 就被推到 500 了。37 fail 了也无法重新消费,offset 已经过去了。
06 DLQ 的 API 极简,但它的 failure surface 是五个维度的
DLQ 的 API 可能只有 produce + consume 两个动作。但 poison pill classification、retry ordering、replay idempotency、depth monitoring、partial-batch ack,每一个都是 production 规模的 distributed systems 问题。出事的时候你才发现这条 backup path 自己也需要 backup plan。



