Cuthbert's Blog

消费者活着、心跳正常、日志干净,但它七天没拉过一条消息

Published on
/1,846 字 · 5 分钟/---

一幕反物理的位点

八月初的一个周一,我在 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已删除...│ ██ 磁盘上实际存在 ██ │
    └──── 真丢 ──────┴──────── 可救 ────────┘

由此得到几条推论:

  1. 消费位点 < 最小位点 = 卡死时长已超 retention 的铁证
  2. 真丢的量 = 最小位点 − 消费位点;还可救的量 = 最大位点 − 最小位点;
  3. 还能反推卡死时刻:生产速率 ≈ 磁盘存量 ÷ 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

  1. 消费位点 < 最小位点,是"卡死已超 retention"的铁证。三位点的判读值得练成条件反射。
  2. 最危险的故障不喊疼。心跳正常 + poll 正常返回空 + 零异常日志,三个"正常"叠加就是七天无人发现。消费 lag 告警不是可选项。
  3. 应用层兜底要覆盖"零异常"形态。我们原有的自愈逻辑只在特定异常时重建 consumer,对这种静音卡死完全无感;补上"连续 N 分钟空 poll 且 lag 在涨 → 重建 consumer"才算闭环。
  4. 客户端库版本是集群级风险面。一个 fetcher bug 在一百多个 pod 里潜伏,审计一次版本的成本远低于下一次事故。
  5. 位点数字本身就是取证材料:生产速率能从"存量 ÷ retention"里算出来,卡死时刻能从堆积量反推出来,多分区独立推算互相印证——不需要日志也能还原时间线。

参考