一幕反物理的位点
八月初的一个周一,我在 Kafka 控制台巡检一个跑了很久的消费组。这个 topic 有 12 个分区,12 个副本的消费服务一一对应,平时安静得没什么存在感。
但那天我在分区列表里看到了一幕怪事:分区 3 和分区 5 的消费位点,比最小位点还小。
如果你对这三个位点的关系有直觉,就会知道这一幕有多反常——正常的消费者是贴着最大位点跑的,无论它多慢,都不该被"最小位点"反超。这就像你在跑道上落后了几圈,结果发现自己落后的那段跑道已经被拆了。
先说结论:两个 pod 因为 kafka-python 2.1.4 的一个 fetcher bug 永久停止了拉取,卡了七天多;期间心跳线程一直正常,所以不触发 rebalance,死死占住分区;poll 每秒正常返回空,日志里零异常。等我发现时,卡死时长已超过 topic 72 小时的消息保留期——约 1157 条消息被 broker 物理删除(真丢),约 814 条还躺在磁盘上等着救,而且这个窗口每小时都在缩小。
下面是完整的排查过程。
三个位点,一把尺子
先把这把尺子立起来,后面所有推理都靠它。Kafka 分区是只追加的日志,retention 只从头部删。控制台给你三列数字:
- 最小位点:磁盘上还存着的最老 offset(log start)
- 最大位点:最新写入的位置(高水位)
- 消费位点:消费组的书签——注意它只是书签,不影响消息的存亡
消费位点 最小位点 最大位点
│ │ │
───┼────────────────┼──────────────────────┤
...retention已删除...│ ██ 磁盘上实际存在 ██ │
└──── 真丢 ──────┴──────── 可救 ────────┘
由此得到几条推论:
- 消费位点 < 最小位点 = 卡死时长已超 retention 的铁证;
- 真丢的量 = 最小位点 − 消费位点;还可救的量 = 最大位点 − 最小位点;
- 还能反推卡死时刻:生产速率 ≈ 磁盘存量 ÷ retention 时长,卡死时长 ≈ 堆积量 ÷ 速率。
第三条后面会用到,而且好用得超出预期。
无异常的空转
从控制台的"客户端 IP"列反查到两个 pod(kubectl get pods -A -o wide 按 IP grep),它们已经 49 天没重启过。看日志,画风是这样的:
每秒一条 No messages received,日志轮转窗口内的 26 个小时里,没有任何异常。
这个画像本身就是重要线索。我们的消费封装里,所有 poll 抛出的异常都会打 Error polling messages——日志里没有它,说明 poll 一直在"正常"返回空结果。这排除了一整类"抛错——重试——再抛错"的问题:它不是在出错,它是在心安理得地空转。
"有异常"的故障好查,日志会喊。"零异常的卡死"是最难的一类——所有观测面都告诉你它很健康。
缺了一条 TCP 连接
接下来需要一个客观事实:这个消费者到底还在不在和 broker 说话?
容器里没有 ss 和 netstat,但有 Python——直接读 /proc/net/tcp:
kubectl -n <ns> exec <pod> -- python -c "
import socket,struct
def ip(h): return socket.inet_ntoa(struct.pack('<I',int(h,16)))
for line in open('/proc/net/tcp').readlines()[1:]:
f=line.split(); l,r,st=f[1],f[2],f[3]
ra,rp=r.split(':')
if st=='01' and int(rp,16)==9092: print(ip(ra))"对比结果一锤定音:健康的 pod 都保持 3 条到 broker 的连接(2 条 coordinator/bootstrap + 1 条到自己分区的 leader);两个卡住的 pod 都恰好缺了到分区 leader 的那条连接。而其它 pod 与同一台 broker 的连接完全正常——broker 是可达的,是客户端自己不再发请求,连接被空闲回收后也不再重建。
嫌疑人范围从"网络、broker、客户端"三方,收敛到了客户端一方。
去源码里找答案
kubectl exec 进容器,pip show kafka-python 确认版本 2.1.4,然后把上游仓库的 fetcher.py 逐 tag diff。答案在 _handle_fetch_response 里:
# kafka-python 2.1.4 _handle_fetch_response
if node_id not in self._session_handlers:
log.error(...); return # ← 泄漏:没有 remove(node_id)
if not self._session_handlers[node_id].handle_response(response):
return # ← 泄漏:fetch session 报错走这里
...
self._nodes_with_pending_fetch_requests.remove(node_id) # 只有成功全程才移除_nodes_with_pending_fetch_requests 集合记录"哪些 broker 上还有在途的 fetch 请求",用来防止重复发。问题是:响应处理有两个提前 return 的分支不清理这个集合。
于是触发链是这样的:broker 侧某次波动 → 返回 fetch session 错误(比如 FETCH_SESSION_ID_NOT_FOUND)→ handle_response 返回 False 提前 return → 这个 node 永远留在 pending 集合里 → 后续 _create_fetch_requests 永远跳过这台 broker——而且跳过时的日志级别是 log.log(0, ...),任何配置都打不出来 → 不再发请求 → 连接被空闲回收,也不再重建。
心跳是独立线程,走的是 coordinator 连接,一直正常——所以消费组不 rebalance,健康副本也无法接管。每一环都"正常",合起来是一台完美的静音故障机器。
Java 客户端历史上有同型 bug:KAFKA-8950。kafka-python 在 2.2.3(2025-05)修掉了它,改成 future.add_both(self._clear_pending_fetch_request, node_id)——成功、失败、session 错误一律清理。我逐 tag 验证:2.1.5 ❌ / 2.2.0 ❌ / 2.2.3 ✅ / 2.3.1 ✅。
它是哪天死的
还有个疑问:卡死是哪天开始的?日志早就轮转掉了,控制台也没有记录。这时候前面那把"位点尺子"就派上用场了:
- 分区 3 磁盘存量 390 条,retention 72h → 生产速率 ≈ 5.4 条/小时;
- 分区 3 总堆积 972 条 → 卡死时长 ≈ 972 ÷ 5.4 ≈ 180 小时 ≈ 7.5 天前;
- 分区 5 同法独立推算 ≈ 7.1 天前。
两个分区各自独立推算,指向同一个时间窗——那几天 broker 侧大概率发生过一次波动(fetch session 被挤出缓存、瞬时网络抖动都可能,云厂商支持查询后答复 broker 侧无维护记录)。触发源最终没有实锤,但这不影响处置:bug 在客户端,升级之后任何扳机都只会正常重建 session。
少丢一点的恢复姿势
恢复不是无脑重启。当消费位点已经越界,新起的 consumer 会发现 offset 不存在,走 auto_offset_reset 逻辑——如果配置是 latest,磁盘上还可救的 814 条也会被直接跳过。
| 方案 | 操作 | 结果 |
|---|---|---|
| 快 | 直接删 pod 重建 | offset 越界 → 跳到最新,可救的 814 条也丢 |
| 少丢 | 副本缩 0 → 控制台重置位点到最早 → 扩回 | 救回磁盘上的 814 条,只丢已被删除的 1157 条 |
我们选了后者,代价是整个 topic 停消费几分钟。这是一个值得记住的决策点:位点越界之后,重启的姿势决定你还能救回多少。
还有多少颗同款地雷
单点修复之外,更重要的问题是:整个集群里还有多少消费者踩在同一个版本上?
版本审计的方法很简单,逐个 pod 执行 python -c "import kafka;print(kafka.__version__)"。结果:同族的 11 个服务、100+ 个 pod 全部是 2.1.4——每一个都可能在下一次 broker 波动时随机引爆。修复因此从"重启两个 pod"升级为"全家批量钉版本到 kafka-python==2.3.1"。
顺手留下一条可复用的卡死巡检命令(对持续 poll 的 kafka-python 消费者):
# 到 broker 的 ESTABLISHED socket 数 < 3 即有嫌疑
kubectl -n <ns> exec <pod> -- python -c "
import socket,struct
n=sum(1 for l in open('/proc/net/tcp').readlines()[1:]
if l.split()[3]=='01' and int(l.split()[2].split(':')[1],16)==9092)
print(n)"Takeaways
- 消费位点 < 最小位点,是"卡死已超 retention"的铁证。三位点的判读值得练成条件反射。
- 最危险的故障不喊疼。心跳正常 + poll 正常返回空 + 零异常日志,三个"正常"叠加就是七天无人发现。消费 lag 告警不是可选项。
- 应用层兜底要覆盖"零异常"形态。我们原有的自愈逻辑只在特定异常时重建 consumer,对这种静音卡死完全无感;补上"连续 N 分钟空 poll 且 lag 在涨 → 重建 consumer"才算闭环。
- 客户端库版本是集群级风险面。一个 fetcher bug 在一百多个 pod 里潜伏,审计一次版本的成本远低于下一次事故。
- 位点数字本身就是取证材料:生产速率能从"存量 ÷ retention"里算出来,卡死时刻能从堆积量反推出来,多分区独立推算互相印证——不需要日志也能还原时间线。
参考
- kafka-python PR #2607:修复 fetch session 错误路径的 pending 泄漏
- kafka-python 2.2.x changelog
- KAFKA-8950 / apache/kafka#7511:Java 客户端同型 bug

