Kafka消费“假活”18小时?监控全绿但消息积压,根因与排查复盘
监控大屏上Kafka三个broker节点的CPU、内存、磁盘指标清一色绿色消费组的状态也显示正常没有任何告警弹出来。但实际情况是一个核心订单topic的消费已经停摆了整整18小时消息从午夜开始越积越多直到第二天上午运营同事跑来问“为什么订单数据一直没更新”我们才意识到出事了。这个案例是典型的“无感故障”Kafka本身没挂、broker没宕机、消费者进程还活着、平台监控也没有触发任何报警但业务链路早就瘸了。今天把这18小时的排查过程完整复盘一遍包括当时怎么定位、为什么监控没有发现、根因出在哪以及事后我们做了哪些改造希望能给正在用Kafka做消息中间件的团队一些参考。1. 事故全景大屏全绿业务却在流血1.1 18小时发生了什么先把现场还原一下。那是一个普通的午夜我们某个核心订单服务的上游系统在凌晨发起了一波批量数据推送单量大概是平日的三倍左右。这个服务通过Kafka接收上游订单消息消费者拿到消息后会调用一个下游的会员积分接口做数据补充。问题出在这个体验积分接口上。凌晨上游系统在做一次数据迁移接口响应从正常的50毫秒直接飙升到几十秒很多请求直接挂起不返回。我们的消费逻辑里面调这个接口时没有设置读取超时时间HTTP连接建立了之后就一直干等。消费者线程陆续全部堵在下游调用上poll循环拉不到新的消息offset也不再提交。从凌晨2点左右开始这个topic的消息积压量就以肉眼可见的速度增长。到第二天上午10点积压的消息已经有几十万条消费完全停滞。但有意思的是这期间没有任何监控告警平台监控大屏上Kafka相关的所有指标都是绿的。1.2 最可怕的不是故障是“无人感知”复盘的时候我们达成了共识这个故障最可怕的地方不在于消费停了18小时而在于18小时里没有任何人发现。下游接口恢复之后积压的消息靠Kafka自身的重试机制缓慢消化整个过程看起来就像什么都没发生过一样但业务侧的数据延迟已经造成了实质影响。这类故障在Kafka使用中相当典型。Kafka的broker节点设计得非常健壮它本身不太容易出问题大量故障其实发生在消费链路——消费者进程还活着但业务处理已经卡死消息处理速率降为0broker端却什么都看不出来。如果监控体系只覆盖了主机层和中间件层没有覆盖到业务消费质量这一层那这种“卡死”就会成为监控体系的盲区。2. 现场排查从“消费不动”到“根因落地”2.1 先用这条命令确认消费组状态接到运营反馈后第一件事是确认消费组和消息积压情况。我们当时的Kafka版本是2.8客户端是Java 11所以直接用了自带的命令行工具kafka-consumer-groups.sh --bootstrap-server kafka01:9092,kafka02:9092,kafka03:9092 \ --group order-center --describe输出大概长这样GROUP TOPIC PARTITION CURRENT-OFFSET LOG-END-OFFSET LAG CONSUMER-ID order-center order-topic 0 12843001 12900312 57311 consumer-1-xxx order-center order-topic 1 11320984 11347520 26536 consumer-2-xxx order-center order-topic 2 22017890 22107903 90013 consumer-3-xxx order-center order-topic 3 9823418 11023345 119927 consumer-4-xxx一眼就看出了问题所有分区的LAG都在几万到十几万之间而且再过两分钟重新执行一次CURRENT-OFFSET几乎没有变化LOG-END-OFFSET却在稳定增长。这意味着生产者一直在往topic里写消息但消费者已经完全不消费了。这里有一个判断技巧如果LAG只高不动要看是CURRENT-OFFSET不动还是LOG-END-OFFSET不动。前者代表消费者卡死后者代表生产者断流。如果LAG在缓慢变小说明消费者还在处理只是速率跟不上如果LAG在持续变大说明消费已经完全停滞或速率远低于生产速率。2.2 区分“消费者掉了”还是“消费者假活”确认LAG异常之后下一步要判断消费者是已经退出消费组了还是“假活着”。再看一遍上面的输出CONSUMER-ID列还有值说明客户端进程还挂在消费组里没有掉线。这里补充一个运维经验如果CONSUMER-ID列是空说明这个分区已经没有活跃消费者在接手通常是消费者进程被broker判定超时踢出了消费组。如果CONSUMER-ID还有值说明消费者的会话还维持着但消费逻辑可能已经卡死。为了确认消费者是否真的还活着我们用jstack抓了一下消费者进程的线程栈jstack 消费者pid jstack_$(date %s).log grep -A 30 kafka-consumer jstack_*.log | head -100线程栈里有大量业务线程停留在HTTP调用的socketRead0方法上状态是WAITING。绕过中间过程结论已经比较明显消费者业务线程全部阻塞在下游接口上没有超时机制一直在等响应。真正的卡死不在Kafka客户端而在消费链路的下游依赖。2.3 线程栈与日志把“卡住”的位置揪出来线程栈只能看到最终状态要想知道“卡了多久”“从什么时候开始的”还得看日志。消费者客户端的日志里有几个关键线索[INFO] Attempt to heartbeat failed since group is rebalancing [INFO] Resetting generation due to heartbeat failure [WARN] Consumer member order-center-xxx has failed, removing it from the group出现这些日志说明消费者业务线程卡死后poll不再被调用超过了max.poll.interval.msbroker判定该消费者失联开始触发再均衡。但由于所有消费者实例都在相同的下游调用上卡住再均衡完成后新当选的消费者依然处理不了消息于是陷入“处理超时—踢出—再均衡—又卡死”的循环。排查到这里根因已经比较明朗了。我当时还看了一个细节进程的CPU占用率并不高但线程数非常多大量线程停留在WAITING状态。这种“进程活着但业务已经不动”的状态比进程直接崩掉更难发现——因为进程活着意味着进程存活层面的监控完全失效。2.4 三个容易误判的细节这个case排查下来有几个细节值得单独说一下都是容易造成误判的地方。第一个是业务日志和系统日志的区分。消费卡死期间业务日志里可能还在持续打印一些上下文日志比如“开始处理消息ID xxx”但永远不会打印“处理完成”。如果你只看日志量会觉得系统还在运转实际上每条消息都卡在中间环节日志量可能还不少。排查的时候一定要看处理完成的日志有没有持续产出。第二个是rebalance日志的迷惑性。出现rebalance日志的时候很多人第一反应是消费者频繁上下线导致的问题会去查网络、查session.timeout配置但忽略了一个前置问题为什么poll停止调用了。rebalance只是结果不是原因原因在业务处理线程的阻塞。第三个是不要只看Kafka侧。Kafka的broker日志、broker指标都很正常因为broker只是存储和分发消息消费者的处理不在它的管辖范围。遇到消费卡死要把排查重心放在消费者进程本身的线程状态和下游依赖上。3. 根因深挖Kafka消费者的“假活”机制3.1 心跳线程和业务线程是两回事很多刚接触Kafka的人会把消费者理解成一个简单的“拉消息—处理—提交offset”循环觉得进程活着就是在消费。但实际上Kafka消费者的客户端内部是分线程工作的。在Kafka 0.10.1之后心跳发送被移到了独立的后台线程中。也就是说即使你的业务处理线程卡在某个下游调用上长时间不执行consumer.poll()心跳线程依然会按heartbeat.interval.ms的间隔向broker发送心跳看起来消费者还是“活着”的。broker侧只关心心跳是否正常它无法感知你的业务线程是否卡死。这套设计本身是为了避免消费者处理超长任务时被误判失活但副作用就是当业务线程真正卡死的时候系统从Kafka broker的视角看一切都是正常的。消费者还在消费组里分区分配也没有变化只有LAG在缓慢增长。这也是为什么这个案例里broker层的监控指标全绿——因为broker确实没有感知到任何异常。要真正识别“假活”只能靠消费延迟类指标。看消费者有没有在处理消息最好的方式就是看CURRENT-OFFSET有没有持续往前移动或者看消费速率是不是趋近于0。3.2 最该背锅的三个客户端参数这个case里有几个Kafka客户端参数是我们事后重点复盘的对象建议所有用Kafka的团队都检查一遍。第一个是max.poll.interval.ms默认值是300000也就是5分钟。这个参数定义了两次poll()之间允许的最大间隔如果业务处理时间超过这个值消费者会被判定为处理超时触发离组和再均衡。很多消费卡死场景变成“雪崩”就是因为这个参数和实际业务处理时间不匹配。处理时间本来就很长的业务应该把这个值调大但不能无限调大因为调大了会降低故障发现的时效性。第二个是max.poll.records默认值是500。单次poll返回的消息条数直接决定了单次循环的处理时长上限。如果每条消息的平均处理时间是200ms500条消息就是100秒超出max.poll.interval.ms默认值只是时间问题。反过来如果把这个值调小比如调到100或者50每次poll的处理窗口缩短即使下游慢单轮处理也不会拖太久留给心跳和rebalance的余量就更大。第三个是session.timeout.ms和heartbeat.interval.ms。这两个参数控制心跳超时检测。如果把session.timeout.ms设得太小网络稍微抖动一下消费者就会被认为失活产生不必要的rebalance设得太大消费者的真实故障又很难被及时发现。新版Kafka默认的session.timeout.ms是45秒heartbeat.interval.ms是3秒这个搭配在大多数场景下是合理的。我们当时的配置问题恰恰出在第二和第三个参数上max.poll.records偏大下游一慢单轮处理时间马上超过max.poll.interval.ms然后触发rebalance。理想状态下这个数字要配到“最坏情况下单轮处理时间”不超过max.poll.interval.ms的60%。3.3 99%的卡死根子在下游超时说实话Kafka消费者自身配置出问题导致卡死的情况并不多绝大多数卡死的根因都在下游依赖。这个case就是一个典型HTTP调用没有设置读取超时上游系统一慢消费者就跟着一起堵死。我用一个生活化的类比来解释Kafka的topic相当于一条高速公路生产者是源源不断的来车消费者是收费站的收费员。正常情况下车来一辆收一辆车道保持畅通。但如果收费员处理一辆车要磨蹭十分钟下游接口慢后面的车就会全部堵在高速上。更糟糕的是高速收费站的系统还认为收费员“在岗”心跳正常所以上面的人完全不会意识到这里已经堵了18小时。下游依赖的卡死形态不止HTTP一种。比如数据库连接池被打满获取连接的线程全部阻塞在等待连接上比如调用了第三方接口对方一直不返回又没有设置超时比如使用的线程池队列是无界的任务全部堆在队列里排队等待执行。这些场景最后的表现都一样消费者线程不干活了但进程还活着。所以排查Kafka消费延迟的时候不要只盯着Kafka的参数调优更重要的是把消费链路上所有外部调用的超时时间梳理一遍。连接超时要设读取超时也要设而且读取超时往往比连接超时更容易被忽略。4. 监控盲区复盘为什么全绿骗过了所有人4.1 我们之前到底监控了什么事故复盘的时候我们第一件事就是梳理监控项把当时已有的Kafka监控全部列了出来列完发现确实很“绿”是有原因的。监控对象监控指标当时的配置是否能发现本次故障Broker节点CPU、内存、磁盘、网络Zabbix主机监控否本次broker完全正常Kafka进程进程存活状态进程探活否消费者进程未退出JMX指标请求处理耗时、网络线程利用率未覆盖否压根没监控消费组消费LAG总量按组平均值阈值10万否单分区最高才12万平均后被抹平消费组消费速率环比无否根本没有指标很明显当时的监控体系是“基础设施视角”的能回答机器挂没挂、进程活没活但回答不了“业务是否在正常运转”这个问题。Kafka中间件的引力在于它是链路中的一环但它本身正常不代表链路正常。消费端的健康度必须通过消费延迟、消费速率这些业务视角的指标来体现。4.2 “平均LAG”把单分区故障抹平了这次监控漏报的主要原因之一就是LAG监控的粒度太粗。我们虽然配置了LAG指标但看的是整个消费组所有分区的平均值。这个topic有几十个分区大部分分区因为流量相对均匀LAG其实并不高。卡死的分区确实有十几万的积压但分摊到所有分区一平均LAG平均值只有几千远低于我们设置的10万告警阈值。告警自然就不会触发。这个教训非常深刻LAG监控如果只按消费组的整体维度去算单分区的卡死很容易被平均掉。正确的做法是至少按topicpartition维度去监控LAG而且告警条件要同时覆盖两种情况——单个分区的LAG超过阈值或者整个消费组的LAG总量超过阈值。两者是“或”的关系不是“与”的关系。类似的思路也适用于消费速率的监控。我们后来把消费速率做成了按partition维度的曲线一旦某个分区的消费速率降为0即使其他分区正常这个异常也会很显眼。4.3 告警疲劳报了等于没报还有一个非常现实的问题告警疲劳。其实在事故复盘的时候我们发现卡死后的第二天早上LAG总量其实已经超过了10万阈值理论上告警应该触发了。但值班群里那个时间段本来就有接近二十条其他告警一条“消费组LAG超过阈值”混在里面压根没人点开看。告警疲劳是监控体系里一个隐蔽又致命的问题。告警太多值班人员就会选择性忽略久而久之所有的告警都变成了背景噪音。这个问题的解法不是加告警渠道、加大声量而是做告警分级和收敛。我们后期把告警分成了三个级别P0级是核心业务消费停止直接电话加短信加群里P1级是消费延迟超过30分钟群里值班人P2级是一般性指标波动只在值班群通知不做强打扰。同时配置了告警去重和恢复通知因为告警的“恢复通知”也很重要可以降低值班人员对告警的焦虑感。5. 监控体系补强从“假绿”到“真绿”5.1 消费延迟监控三个维度缺一不可经过这次事故我们重新梳理了Kafka消费链路监控的指标体系核心是三个维度消费组维度、topicpartition维度、消费速率趋势维度。消费组维度解决的是“整个业务链路是否健康”的问题。每天业务量有波动LAG的绝对值会跟着波动所以消费组维度的告警阈值不能只看绝对值更要看变化趋势。比如消费延迟在持续增长就说明消费速率跟不上生产速率即使LAG绝对值还没到阈值也已经是个危险信号。topicpartition维度解决的是“局部是否健康”的问题。Kafka的分区机制决定了单个分区卡死是常态场景所以必须能看到每个分区的LAG曲线。如果某个分区的LAG持续走高而其他分区正常那就是典型的单分区故障需要尽快定位为什么消费者的消息处理在这个分区上卡住了。消费速率趋势维度解决的是“消费是否停滞”的问题。LAG是存量指标消费速率是流量指标两者要结合看。消费速率突然降为0哪怕LAG还很小也代表着消费者可能已经卡死了这个告警的时效性比LAG绝对值告警高得多。5.2 Prometheus kafka_exporter 的落地示例工具选型上我们后来统一用的是Prometheus加kafka_exporter配合Grafana做面板展示夜莺那边也做了指标汇总和告警推送。kafka_exporter是开源的Kafka指标采集器默认会导出大量JMX指标其中消费组相关的核心指标是这两个kafka_consumergroup_current_offset kafka_consumergroup_lag建议在Grafana里单独建一个“消费组健康度”的面板至少要包含这几个视图消费组列表、各消费组LAG总量、各topic下分区的LAG分布、消费速率变化曲线。消费速率可以间接计算通过increase函数看一段时间内offset的变化量sum(increase(kafka_consumergroup_current_offset[5m])) by (consumergroup, topic)这条PromQL的意思是计算过去5分钟每个消费组在每个topic上的offset增量也就是这5分钟实际消费了多少条消息。如果这个值持续为0说明消费停滞如果这个值长期低于生产速率说明消费能力不足。我们后来给这个指标配了一条告警规则连续15分钟消费增量小于某个基线值就触发P1告警。Prometheus部署本身不复杂关键是把exporter的采集粒度调好特别是topics和consumer groups的过滤不要让exporter去采集那些废弃topic和废弃消费组否则指标数量会非常庞大浪费存储资源。5.3 Burrow 和专门的消费组监控工具除了自己搭配的Prometheus方案还想提一下Burrow它的消费延迟评估思路非常值得参考。Burrow是专门做Kafka消费组监控的工具它引入了“窗口评估”的概念不会简单地看LAG绝对值而是把LAG的历史趋势拉出来结合消费速率和分区分配情况判断消费者的状态是否健康。Burrow会把消费组状态分为OK、WARNING、ERROR、STOPPED几种。STOPPED表示消费组已经完全不消费这正好对应我们这次事故的状态。如果当时用了Burrow卡死几分钟内就会被标记出来根本不可能等18小时。不过Burrow的部署和维护成本偏高对Kafka版本的兼容性也需要验证。如果团队规模不大用Prometheus加kafka_exporter做基础监控再配好告警规则已经能覆盖大多数场景。Burrow更适合Kafka集群规模较大、消费组数量很多的公司。5.4 告警分级与值班联动监控指标配好了还不够告警链路必须跑通。我们这次事故的教训之一就是光有告警不行告警必须能触达真正能处理问题的人。现在的配置逻辑是这样的Prometheus负责采集和评估指标Alertmanager负责把告警转发到企业内部的消息平台和电话平台再按告警级别走不同的触达策略。P0级告警会直接拨打电话P1级告警在群里值班人P2级告警只记录到值班群。同时设置告警静默窗口同一个告警规则在2小时内不重复触发避免刷屏。这里有一个建议告警恢复通知一定要配。告警恢复了值班人员才能确认问题已经闭环否则告警虽然触发了但什么时候恢复、有没有恢复完全看不到下一班的人接手会很被动。6. 代码和配置层面的免疫力6.1 给所有下游调用加上超时上限监控只能解决“及时发现”的问题真正让系统不容易出事的还是代码层面的防御。这个case之后我们做的第一件事是全面排查消费者代码里的外部调用凡是HTTP调用必须显式设置连接超时和读取超时凡是数据库调用必须设置连接获取超时和查询超时。以前很多团队写HTTP调用不设置超时是觉得反正内网调用足够快超时是多余的。但内网服务一样会挂上游一样会因为各种原因变慢。设置超时不会影响正常调用但在异常情况下它是消费者卡死前的最后一道防线。连接超时和读取超时要区分开连接超时是建立连接的最长等待时间读取超时是发完请求到拿到响应的最长等待时间。很多代码只设置了连接超时没设置读取超时结果连接建立了服务端那边不返回数据客户端就永远等下去这恰好就是这次事故的形态。6.2 poll循环的正确写法消费端代码的结构也很关键。参考我们修复后的写法while (true) { ConsumerRecordsString, String records consumer.poll(Duration.ofMillis(500)); if (records.isEmpty()) { continue; } // 把业务处理和poll循环分离避免业务处理拖住poll调用 processRecords(records); consumer.commitSync(); }这里有几个要点。第一poll的超时时间不要设置太长500毫秒到1秒比较合适这样可以及时响应rebalance请求。第二业务处理逻辑要控制在max.poll.interval.ms之内如果单批消息处理时间可能很长要么调小max.poll.records要么把处理逻辑放到独立的线程池里执行poll线程只负责拉取和提交offset。第三commitSync和commitAsync要选对处理完一批消息后建议用同步提交确保offset不会丢如果对吞吐要求很高可以用异步提交加回调处理失败情况。另外一种更稳健的模式是使用分离的线程池来处理消息poll线程把消息提交给业务线程池后就立即返回再配合手动控制提交频率。这种模式的复杂度更高但能有效隔离单条消息卡死对poll循环的拖累适合对稳定性要求高的核心链路。6.3 死信队列与重试边界这次事故的第二个直接原因是消费者在消息处理失败时采用了无限重试的策略。下游接口挂起之后每条消息都在那里反复重试重试本身又占用了处理线程导致后面的消息全部排队等待。正确的姿势是给重试设置边界。单条消息的重试次数要有限制超过重试次数应该转入死信队列DLQ。死信队列里的消息可以后续单独处理但不应该阻塞主消费链路。实现上有两种常见方案一种是用单独的topic做死信队列消费者处理失败的消息序列化后写到DLQ另一种是在业务表里加状态字段处理失败的消息记录到一个失败表用定时任务去扫描重放。重试本身也要考虑退避策略。固定间隔的重试在高并发场景下会形成重试风暴建议用指数退避加随机抖动。比如第一次重试等1秒第二次等2秒第三次等4秒加上不超过500毫秒的随机偏移这样既能给下游恢复留时间又不会瞬间打爆下游服务。6.4 消费组的生命周期管理最后提一个很基础但很多人忽略的点消费组的生命周期管理。我们当时Kafka集群里积压了一堆废弃的topic和对应的遗留消费组有的是项目下线了没清理有的是测试环境误建了生产topic的消费组。这些废弃消费组带来的问题是如果你想做消费组维度的监控它们的数量会淹没真正需要关注的消费组。topic一多LAG的告警规则要么就只能针对核心topic单独配置要么就选择全部监控然后被废弃topic的噪音干扰。我们后来做了一次大清理下线了所有废弃topic删掉了所有不再使用的消费组然后对核心消费组做了白名单式的监控配置告警准确率一下子提高了很多。消费组的创建规范也值得定一下比如消费组命名必须包含业务线和应用名这样通过名字就能认出这个消费组是干什么的、该由哪个团队负责。规范定好之后责任归属清晰告警发出来大家也能快速判断该找谁。这次事故之后我养成了一个习惯每周至少要把核心消费组的lag曲线翻一遍不看平均值看每个分区的曲线。有时候最危险的不是没有监控而是监控给了你一种“一切正常”的错觉。Kafka这个中间件本身很稳定大量故障都发生在消费端而消费端是否健康最终要由消费延迟这类业务指标来回答。如果你的监控体系里还没有消费延迟的视角建议尽快补上别等一个topic堵了18小时再来复盘。
上一篇/下一篇内容由系统自动关联
返回资讯列表 →