ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

Kafka消费者假死18小时无人发现?监控盲区与排查实战

Kafka消费者假死18小时无人发现?监控盲区与排查实战 先说结论这次故障最可怕的不是 Kafka 消费卡死本身而是它整整堵了 18 小时所有监控指标看上去都非常正常。如果不是业务方主动反馈“消息发出去之后一直没有回执”这条 topic 的积压事故可能还会继续绿下去。当时接到电话是凌晨 1 点 47 分第一反应先看一眼监控大盘broker 的 CPU、网络吞吐、请求队列、分区在线状态全是绿的。消费者实例的 JVM 堆内存也在正常水位GC 频率没有异常。但切到消费组页面CURRENT-OFFSET已经固定在一个数字上将近 18 个小时LAG曲线从昨天上午开始一路拉升相当于把过去一整天的数据全部积压在了这条 topic 里。那一刻我就明白了这不是 Kafka 集群的事而是消费者侧“假死”了并且我们的监控体系对这种情况完全没有感知能力。这篇文章把整个排查过程完整拆开包括为什么监控会全绿、怎么从消费组状态一步步挖到线程栈、最终定位到哪一行代码、以及后续我做了哪些改造让“消费卡死”能被自动发现。适合所有自己维护 Kafka 生产集群或者正在用 Kafka 做异步消息处理的后端同学参考。1. 故障现场复盘topic 停了 18 小时监控面板上却一片绿色1.1 故障时间线从第一条积压消息到被业务方发现事后拉取时间线问题其实很早就埋下了前一天上午 09:42某条业务消息触发了消费者内部一个外部接口调用这个调用之后一直没有返回。09:47消费者线程还在等待外部响应poll()不再被调用消费行为实际上已经停止。10:00 到 17:00该 topic 的日志终端偏移量LOG-END-OFFSET持续增长但消费偏移量CURRENT-OFFSET纹丝不动。这里开始出现大量的消息积压。17:30 左右下游另一个系统出现了一批“任务超时”的异常但因为重试机制的存在业务还能勉强维持。第二天凌晨积压量已经涨到几百万条业务方终于发现“消息出去后怎么半天没反应”才通过值班渠道找到我们。有意思的地方在于积压并不是瞬间发生的而是像温水煮青蛙一样慢慢累积。如果消费卡死发生在 10 分钟以内监控上可能完全看不出异常一旦超过一个小时LAG指标一定会出现明显抬头。问题是当时的监控体系里压根就没有针对LAG的告警规则所以它涨得再高也没有任何人收到通知。1.2 线上系统的部署形态一个消费者实例单扛所有分区这个出问题的 topic 有 12 个分区但消费者组里只部署了 1 个应用实例。不是说不能这么用Kafka 允许一个消费者进程消费多个分区但这种部署形态有个致命弱点单实例一旦卡死整个 topic 就处于“无人消费”的状态没有其他实例可以接管。事后复盘时团队里有人提出“为什么不做成多实例”这个方向没问题但多实例并不能解决“卡死”本身。即使有 3 个消费者实例如果它们内部都同步调用同一个下游接口下游一挂3 个实例会一起卡住。消费者的高可用只能解决“进程宕机”这类硬故障解决不了“逻辑卡死”这类软故障。这也是这次故障给我们上的最重要一课。2. 为什么 18 小时没人发现监控体系的三层盲区2.1 第一层盲区监控采集与告警规则的错位出问题之后我第一件事就是去翻监控系统。我们当时在 Grafana 上给 Kafka 配了十几个面板broker 吞吐量、请求平均延迟、分区在线状态、磁盘使用率、消费组状态看似应有尽有。但这些面板背后的告警规则呢几乎没有。具体来说broker 层的告警规则是 CPU 超过 85% 或磁盘使用率超过 80%这次 broker 自己非常健康完全没触发。消费组相关的面板有展示LAG曲线但从来没有配置过告警。更关键的是告警阈值怎么设如果拍脑袋把阈值定成 10 万条这次积压到 300 万条才可能被触发到那时已经晚了 10 多个小时。我后来想通了一个道理监控面板是给人看的告警规则才是真正在干活的。面板再全没有对应的告警规则和没有监控没有本质区别。真正的消费监控不应该只关注“是否在消费”而应该关注“是否在及时消费”。2.2 第二层盲区消费者“在线”不等于“在消费”下面这句话是这次事故给我最深的教训Kafka 消费者进程活着和它在正常工作完全是两码事。当时 Kafka 集群的视角里消费者的会话Session是存在的它和 broker 之间的 TCP 长连接也是正常的消费者组的group.state显示为Stable。如果只看 broker 端和消费者端的进程状态一切都正常。但实际发生的情况是消费者的业务处理线程卡死在一个网络读取操作上既没有在做消息处理也没有在调用poll()。进程没有退出连接没有断开线程也没有抛异常。这种情况下Kafka 集群本身根本感知不到这个消费者“工作异常”它只会觉得这个消费者“还活着”。这种假死比进程宕机可怕得多。进程宕机至少会触发连接断开和 rebalance有经验的运维一看消费组成员变化就能发现问题。而假死是静默的如果不关注消费偏移量你甚至不知道它已经停止了。2.3 第三层盲区指标有了但缺少业务闭环验证再往深一层说我们的监控体系还缺了一个“业务闭环”的验证手段。当时我们虽然采集了LAG、消费速率这类指标但这些指标只能告诉你“消费慢了”或者“堆积了”没法告诉你“消费出来的消息业务是否真的处理成功了”。一条消息从生产到被 Kafka 拉取再到业务代码真正处理完成这中间还有很长一段路。如果消息被拉取后业务处理一直不返回Kafka 层面是看不到任何异常的。所以后来我们建设了一套端到端的延迟监控消息生产时在消息体里写入一个时间戳业务消费完成后把“当前时间 - 消息生产时间”作为延迟指标暴露出去。只有这个指标才能真实反映业务链路是否健康。监控的最终目标不是看 Kafka 有没有消息堆积而是看业务有没有被及时处理。3. 从 GROUP 状态到线程栈一步步把“卡死”位置抠出来3.1 第一步kafka-consumer-groups.sh看消费停滞接到告警后我先用命令行看消费组的状态这是排查这类问题最快的方式kafka-consumer-groups.sh \ --bootstrap-server 10.0.0.10:9092 \ --describe \ --group order-delay-consumer输出大概是下面这个样子GROUP TOPIC PARTITION CURRENT-OFFSET LOG-END-OFFSET LAG CONSUMER-ID HOST order-delay-consumer order-delay-queue 0 182345 420156 237811 consumer-1-xxxxx /10.0.0.8 order-delay-consumer order-delay-queue 1 165732 398765 233033 consumer-1-xxxxx /10.0.0.8 ...这里最刺眼的就是LAG列12 个分区的积压加起来已经超过 300 万条。而CONSUMER-ID那一列只有同一个消费者实例说明 group 里没有其他活跃消费者可以分担。我在排查时会重点比较CURRENT-OFFSET的变化趋势连续执行两次命令如果CURRENT-OFFSET完全没有变化而LOG-END-OFFSET一直在增长就能断定这个消费者已经停止消费了。3.2 第二步jstack抓线程栈找到钉死在 IO 上的线程确认消费者停滞之后接下来要回答的问题是它到底卡在哪里因为 Kafka 消费者本身是纯 Java 的所以直接抓 JVM 线程栈是最直接的手段。先找到消费者进程的 PIDjps -l | grep consumer然后连续抓几次线程栈jstack pid /tmp/thread_dump_1.txt sleep 5 jstack pid /tmp/thread_dump_2.txt对比两份线程栈重点看同一个线程是否多次停留在同一个栈帧上。如果连续两次都在同一个位置说明这个线程已经卡了很久。当时抓到的最可疑线程栈局部是这样的consumer-1 #29 prio5 os_prio0 tid0x00007f8c9800 nid0x1234 runnable [0x00007f8c9000] java.lang.Thread.State: RUNNABLE at java.net.SocketInputStream.socketRead0(Native Method) at java.net.SocketInputStream.read(SocketInputStream.java:150) at java.net.SocketInputStream.read(SocketInputStream.java:121) at org.apache.http.impl.io.SessionInputBufferImpl.streamRead(SessionInputBufferImpl.java:137) at org.apache.http.impl.io.SessionInputBufferImpl.fillBuffer(SessionInputBufferImpl.java:153) ... at com.example.orderdelay.consumer.MarketingClient.send(MarketingClient.java:88) at com.example.orderdelay.consumer.OrderDelayProcessor.process(OrderDelayProcessor.java:37)注意这个线程状态是RUNNABLE不是BLOCKED也不是WAITING。SocketInputStream.socketRead0是 Java 中网络读取的系统调用在等待远端数据时线程状态依然是RUNNABLE但实际上它就是在干等。这也是线程栈排查里特别容易迷惑新人的地方不是只有BLOCKED和WAITING才算卡死长时间停滞在 IO 读上同样是卡死。连续抓了 5 次线程栈每一次MarketingClient.send()下面的方法栈都一模一样可以确定这个线程已经卡在这里超过 5 分钟了。3.3 第三步连接池与下游依赖的“连带堵车”定位到具体线程还不够我还想搞清楚一件事是不是只有这一个线程卡住了用jstack输出的线程总数统计了一下消费者服务一共有 40 多个线程但真正在执行 Kafka 消息处理的只有 3 个当时用了一个简单的多线程消费封装这 3 个线程全部停在SocketInputStream.socketRead0上。这意味着不是单条消息的偶发问题而是所有处理线程都堵在了同一个外部接口上。这时再去看依赖的下游接口状态很明显第三方营销中台服务从当天上午开始就处于半死不活的状态TCP 连接能建立但 HTTP 响应迟迟不发出来。由于我们的 HTTP 客户端没有设置合理的读取超时线程就只能无限期地等下去。这里还要注意一个连带效应如果消费者内部还用了数据库连接池或线程池外部接口一堵这些资源很快也会被占满。比如数据库连接池里有 20 个连接其中 15 个被“等待外部接口”的线程持有真正需要数据库连接的消息反而拿不到连接最终导致整个消费者服务所有请求全部排队。3.4 顺便排除的干扰项GC 停顿、OOM、CPU 飙高排查过程中我们还做了几项排除这里也列出来供参考避免大家走弯路GC 停顿用jstat -gcutil pid 1000观察了几分钟Full GC频率很低堆内存使用率维持在 40% 左右基本排除 GC 导致的长暂停。OOMjstat和jmap -heap看下来没有内存溢出迹象OutOfMemoryError也没有出现在日志中。CPU 飙高top -Hp pid看线程 CPU 占用整体 CPU 使用率不到 10%。真正 CPU 飙高的卡死通常会伴随线程持续消耗 CPU比如死循环和这次的症状完全不同。内存、CPU、GC 三个指标都正常反而更坚定了我的判断这是一个典型的“外部 IO 阻塞型卡死”不是资源耗尽型故障。4. 根因还原poll 循环里那行“看起来没问题”的同步调用4.1 问题代码一条消息拖死整个消费线程最终定位到的问题代码非常典型简化后的逻辑大概是这样的while (true) { // 这里设置了 1 秒的 poll 超时Kafka 有消息时几乎立即返回 ConsumerRecordsString, String records consumer.poll(Duration.ofMillis(1000)); if (records.isEmpty()) { continue; } for (ConsumerRecordString, String record : records) { // processMessage 内部会同步调用外部营销中台的 HTTP 接口 processMessage(record.value()); } }而processMessage内部真正卡住的是一个没有设置读取超时的 HTTP 调用private void processMessage(String message) { // 发短信通知用户 marketingClient.sendMessage(message); // 后续还有一些数据库操作 }问题就出在marketingClient使用的 HTTP 连接没有配置 socket 读取超时。很多 HTTP 客户端的默认行为是“无限期等待响应”一旦下游服务异常但不主动断开连接调用方就会一直挂在socketRead0上。更麻烦的是processMessage是在poll()返回之后同步执行的它卡住消费者就没办法在预期时间内发起下一次poll()整个消费行为完全停止。这里想特别强调不要在 Kafka 消费的主线程里同步调用不可控的外部依赖尤其是没有超时的调用。如果真的要调用也要用异步的方式或者至少保证外部调用有明确的超时时间、有熔断机制。4.2 Kafka 消费者心跳机制在卡死场景下的真实表现有人可能会问Kafka 消费者不是有心跳机制吗卡死 18 个小时为什么没有被踢出组这个问题的答案取决于卡死的方式以及客户端版本的行为。在 Kafka Java 客户端较老版本的设计中心跳的发送是在poll()方法内部完成的也就是说如果业务线程卡在poll()返回之后长时间不调用下一次poll()心跳也会跟着停。超过session.timeout.ms默认 45 秒后broker 会认为该消费者已经下线触发 rebalance。但在我们的实际情况里有两个因素导致了这个机制没有生效团队里有人为了处理大批量消息时避免误触发 rebalance把max.poll.interval.ms调成了 30 分钟又把session.timeout.ms调到了 120 秒。于是第一次卡死出现后的前半小时内消费者都不会被判死。更重要的是即使后面的某个时间点触发了 rebalance这个消费组里也没有第二个消费者实例可以接管这些分区。Kafka 的 rebalance 只能在“有消费者可以接管”的情况下才能恢复消费group 里只有一台实例它被移除后这些分区就变成无人消费状态只是从 Kafka 的视角来看group 不再稳定而已。所以在单实例消费的场景下心跳超时和 rebalance 机制并不能提供有效的保护。它们能够感知“消费者挂了”但无法感知“消费者还活着但不干活了”。4.3 为什么 18 小时都没触发 rebalance 和 failover把上面两点串起来就能解释这个“18 小时无人发现”的完整链条了消费者的处理线程卡在了外部 HTTP 调用上没有抛异常进程也没有退出。外部服务的 TCP 连接还保持着JVM 和 broker 之间的长连接也没有断开。由于max.poll.interval.ms调得比较大短时间内没有触发消费者离组。即使之后离组了group 内没有其他实例分区依然无人消费消息继续积压。所有现成的监控指标里唯一会变化的就是LAG但我们没有针对它配置告警规则。最终业务方通过自己的链路发现问题时已经过去了 18 小时。这串链条里每一个环节单独看都不致命但合在一起就成了一次长时间的生产事故。排查 Kafka 消费故障时一定要想清楚你的监控到底能不能捕捉到“消费逻辑假死”这类故障5.1 紧急恢复跳过积压消息和临时迁移分区的操作记录紧急恢复阶段我们没有马上让消费者去追 300 万条积压消息。那样做只会让它继续被历史消息里的外部调用卡住治标不治本。实际做了三步操作第一步确认消费者进程安全直接重启。这一步的目的是让所有卡死的线程重新初始化释放被占用的 TCP 连接和线程池资源。第二步跳过积压消息先保证新消息能及时消费。重启后如果消费者继续从旧的 offset 开始拉取它第一时间还是会读到那条“让外部调用卡死”的消息可能再次卡住。所以需要把消费位置重置到最新kafka-consumer-groups.sh \ --bootstrap-server 10.0.0.10:9092 \ --group order-delay-consumer \ --topic order-delay-queue \ --reset-offsets --to-latest --execute注意--to-latest这个操作会把当前 group 的 offset 全部移动到最新位置意味着一部分积压消息会被跳过。如果业务不能丢消息需要先把积压数据导出或者用另一个临时 group 慢慢追而不是直接重置。我们的业务场景允许在恢复期跳过存量消息所以这一步是安全的。第三步在代码修复完成之前先临时摘除消费者对异常外部接口的依赖。具体做法是在代码里加了一个开关默认关闭外部营销中台调用保证消息能正常消费。下游恢复后再逐步打开开关。这三步做完topic 的消费恢复到了正常状态新的消息延迟降到了秒级。积压的 300 万条消息我们后续用了另一个临时消费者慢慢补处理。5.2 长期改造消费逻辑的超时、隔离与异步化问题恢复之后代码层面的修正是必须的。我不会只改一行超时时间就完事而是把消费者侧的工程化水平系统性地提升一档。首先是给所有外部调用加上明确的超时时间。HTTP 客户端连接超时、读取超时都要设置连接超时建议 3 秒以内读取超时可以按业务场景设 5 到 10 秒绝对不能是无限大。这一步虽然基础但至少能保证单条消息的阻塞不会无限期拖住整个消费线程。其次是把消费处理改成异步或独立线程池模式。不要在poll()返回后的主循环里同步执行耗时的外部调用。改造后的模型大致是这样的// 单独线程池处理业务逻辑 ExecutorService bizExecutor new ThreadPoolExecutor( 4, 8, 60, TimeUnit.SECONDS, new ArrayBlockingQueue(2000), new ThreadPoolExecutor.CallerRunsPolicy() ); while (true) { ConsumerRecordsString, String records consumer.poll(Duration.ofMillis(1000)); for (ConsumerRecordString, String record : records) { bizExecutor.submit(() - processMessage(record.value())); } // 手动提交offset或在业务处理成功后提交根据业务可靠性要求选择 }这种设计的核心价值是poll()循环永远不会被单个消息的处理时长卡住即使某条消息处理得再慢也只是占用了线程池里的一个线程后续消息还能继续被拉取和处理。但需要注意改成异步之后原来“一条消息处理完再消费下一条”的顺序语义会被打破必须评估业务是否接受乱序。最后是给消费者加上熔断和降级机制。当外部依赖连续失败达到阈值时直接快速失败而不是继续死等让消息进入死信队列或本地失败表后续人工补处理。5.3 监控补齐真正能抓住“消费卡死”的三类告警线上的代码修复是一方面监控体系的补齐才是保证下次不再发生“18 小时无人发现”的关键。我梳理了以下三类必须配置的告警第一类消费延迟Lag持续增长告警。单纯给LAG设一个固定阈值其实不太靠谱因为不同 topic 的积压容忍度完全不同。更合理的策略是监控“Lag 是否持续增长”比如每隔 5 分钟采集一次如果连续 6 次都发现 Lag 在涨就触发告警。这种基于趋势的告警比固定阈值敏感得多也更能反映“消费者停滞”的特征。第二类消费者组活跃成员数量告警。在 Kafka 中如果消费者实例真的宕机了group 的活跃成员数会减少。给这个指标加上告警可以在一定程度上弥补“在线但假死”的盲区。虽然它抓不到“实例活着但业务卡死”的场景但至少能让进程级异常第一时间暴露。第三类端到端消息处理延迟告警也就是探针消息方案。这一类是我认为最有效、也最值得推荐的。基本思路是准备一个专门的探测 topic定期比如每 5 分钟往里发一条带有时间戳的“探针消息”消费者启动后自动订阅这个 topic消费到探针消息后把“当前时间 - 探针消息时间戳”上报给监控系统。如果延迟超过设定的阈值比如 15 分钟就说明消费者的业务处理链路已经出了问题。这个方案的厉害之处在于它检测的是“消费者是否真的在处理消息”而不是“消费者进程是否存在”。无论消费者是卡死、死锁、还是被外部依赖阻塞只要它没有及时消费探针消息监控就能感知到。我们后来在线上所有核心消费组都部署了这套探针机制几乎没有再出现过“消费者挂了但没人知道”的情况。除了上面三类告警我们还在消费者代码里主动暴露了一批自定义指标比如last_poll_timestamp、last_processed_timestamp、业务线程池队列深度等接入 Prometheus 后可以直接用 PromQL 查询消费者“最近一次 poll 是什么时候”“线程池里积压了多少任务”。这些指标在下次排查类似问题时能帮你在分钟级定位到问题范围。这次故障之后我最大的感受是Kafka 的监控很容易做得“看起来非常完善”但真正能扛住事故的往往是那些最朴素的东西——有告警规则的指标、有业务语义的探针、以及消费者侧对不可控依赖的防护。如果你现在维护的消费组还没有探针消息我建议真的认真考虑加上一套成本很低回报是关键时刻能把 18 小时缩短成几分钟。
返回列表