ARTICLE DETAIL

资讯详情

深耕网站视觉设计与运营推广的一线实战洞察。

Kafka消费积压排查实战:从Lag告警到线程Dump的完整定位指南

Kafka消费积压排查实战:从Lag告警到线程Dump的完整定位指南 凌晨两点十七分告警群里跳出一条消息Kafka消费组order-payment-sync的Lag值突破了五万而且还在持续向上走。我打开监控面板看了一眼消费速率几乎归零生产速率却纹丝不动。那一刻我就知道这个Bug排查的夜晚又要占掉大半宿。这类问题最折磨人的地方在于进程没崩溃日志没报错消费者还活着但业务就是卡住了。对后端工程师来说消费积压、线程阻塞、偶发故障几乎是职业生涯绕不开的三座大山。这篇日记我想完整记录一次从告警到定位的排查过程再把这些年在不同场景里踩过的坑——从Pulsar Key_Shared模式不消费到Java容器内存飙高再到IP冲突、CAN通信和前端断点调试——一起梳理成一套真正能复用的排查思路。不管是刚入门的新人还是已经被线上问题磨过几年的老手这套方法论应该都能帮你少走一点弯路。1. 告警夜里的第一现场Kafka Lag持续攀升消费者却在摸鱼1.1 先别急着看代码确认现象、影响面与时间线收到Lag告警后大多数人第一反应是打开消费代码眼睛在逻辑里来回扫。我的建议是忍一忍。代码不会跑掉但现场会消失。如果连现象、时间线、影响范围都没搞清楚就动手很容易在错误方向上浪费半小时。先看监控大盘把这几项拉出来Lag曲线、消费速率records-consumed-rate、生产速率records-produced-rate、GC暂停时间、网络流量。重点确认三件事一是Lag从什么时间点开始上涨这个时间点附近有没有发布、变更、压测或下游依赖抖动二是积压是全部分区还是部分分区是所有消费实例还是单个实例三是消费速率是直接归零还是缓慢下降。这里有一个基础结论需要先背下来Lag上涨只有三种可能生产速率上升、消费速率下降、消费者彻底不消费。在绝大多数告警场景里问题都出在消费侧。所以接下来要回答的问题是——消费者到底在忙什么。我用一张表把信息收集的动作固定下来每次排查都照着填检查项命令/入口关注点Lag曲线Grafana/Kafka监控起始时间、增长速度、涉及的分区消费速率消费者指标归零还是下降集群状态Kafka UI/Broker日志有无分区Leader切换、磁盘异常消费者组详情kafka-consumer-groups.sh --describeoffset是否停住、rebalance次数应用日志最近30分钟WARN/ERROR有无异常、超时、OOM线程状态jstack采样三次BLOCKED/WAITING线程基础指标CPU/内存/磁盘/网络排除外部资源问题这一套动作做完问题范围基本能缩小一半至少能判断是Kafka集群抖动还是消费端卡住了。1.2 从Lag曲线和消费者组状态反推故障层确认完现象之后马上执行一条命令kafka-consumer-groups.sh --bootstrap-server localhost:9092 \ --describe --group order-payment-sync这条命令会列出消费组里每个分区的当前offset、LogEndOffset、Lag和Consumer实例。重点看两个细节第一所有分区的offset是否都停在了同一时间点如果全部停住基本是消费进程级的问题第二有没有分区在不停发生rebalance如果Member数量频繁变化、JoinGroup事件暴增那可能是心跳线程出问题或者消费耗时太长导致被踢出。我有一次踩过典型的rebalance陷阱业务代码里在poll循环中执行了一个长达五分钟的批量任务session.timeout.ms还是默认的45秒消费者被Broker判定为死亡反复触发rebalance。表面上看Lag在涨、消费者在报错实际根因是消费线程被长任务占住根本没空poll。所以在看消费者组状态时一定要把是否频繁rebalance当成一个重要信号而不是只在日志里找异常。另一个关键动作对比消费速率指标。Kafka的消费者指标里有records-consumed-rate和records-lag-max如果records-consumed-rate已经归零说明poll调用压根没有返回数据问题在应用层面不在Kafka。1.3 第一时间该收集哪些现场数据很多人排错失败不是分析能力不行而是证据链断了。等你想起来要对比数据的时候现场已经被重启或者被新的日志冲掉了。所以一旦判断故障层级在应用马上把下面这些资料收集齐按时间线归档故障发生前后30分钟的应用日志重点找ERROR、WARN、OOM、Timeout关键字连续三次线程Dump间隔3到5秒命令是jstack -l pid内存和GC日志用jstat -gcutil pid 1000看趋势如果怀疑内存泄漏再考虑Heap Dump基础设施指标CPU用户态/内核态、磁盘IO、网络重传率、TCP连接数最近的发布记录、配置变更记录、下游接口耗时。线程Dump尤其重要它可以清楚地告诉你每个消费者线程在某一瞬间具体卡在哪一行代码、哪个锁、哪个方法上。连续采几次还能看出线程是一直卡着还是偶发卡顿。记住一个原则Dump是静态快照多采几次才有对比价值。2. 层层下钻从线程Dump到业务代码的定位全过程2.1 先排除外部依赖网络、磁盘、GC在分析代码之前先把外部因素排掉这一步花不了几分钟但能避免你当背锅侠。第一步看CPU。如果CPU跑满消费者线程可能在死循环或者大量计算如果CPU很低那大概率是线程在等待——等锁、等IO、等网络、等连接池。第二步看GC。执行jstat -gcutil关注Full GC次数和耗时。G1或CMS的Full GC会造成整个应用暂停暂停期间消费线程当然不工作Lag就会快速上涨。我见过一个案例某服务堆大小设成4G但元空间没有上限长时间运行后元空间膨胀频繁Full GC消费线程几乎处于半瘫痪状态看起来像消息堆积问题实际上是被GC拖垮的。第三步看网络和磁盘。消费者到Kafka broker之间的网络重传率异常升高会导致拉取请求变慢应用日志落盘太频繁把磁盘IO打满也会拖慢一切操作。这几项排查完你大概能判断是外部环境引起还是应用自身逻辑引起接下来才轮到线程Dump登场。2.2 线程Dump与锁分析找到一动不动的消费线程线程Dump是定位Java应用卡顿问题的首选工具没有之一。连续执行三次jstack -l pid jstack_1.txt sleep 4 jstack -l pid jstack_2.txt sleep 4 jstack -l pid jstack_3.txt然后找状态为BLOCKED和WAITING的线程。很多情况下你会看到类似的输出kafka-consumer-1 #23 prio5 os_prio0 tid0x... java.lang.Thread.State: WAITING (parking) at sun.misc.Unsafe.park(Native Method) at java.util.concurrent.locks.LockSupport.park(...) at java.util.concurrent.ThreadPoolExecutor.getTask(...) at java.util.concurrent.ThreadPoolExecutor.runWorker(...)如果消费线程在等线程池任务说明提交给线程池的任务一直得不到执行线程池队列满了或者核心线程被占满。如果看到waiting for monitor lock说明锁被别的线程持有你要顺着堆栈找到持锁线程看它在干什么。之前遇到过的情况是某个消费线程在处理消息时调用了一个报表接口该接口内部用了全局同步锁相当于整个消费者被一个慢接口串行化。消费者组二十个线程全部堵在同一把锁上。看线程Dump的时候不要只看Java层还要留意Native层。有时候堆栈停在SocketInputStream.read上说明是在等待网络返回可能是下游接口超时也可能是TCP连接已经半开。结合netstat看一下连接状态可以进一步确认是不是下游依赖导致。回到这次Kafka案例线程Dump显示消费线程在等业务线程池而业务线程池里的worker线程几乎全部卡在SocketInputStream.read上。再配合netstat能看到到下游某个服务的连接数已经占满明显是下游请求被挂在那里了。2.3 真相往往藏在超时和异常处理里在上面这个案例里最终的根因其实很朴素业务代码调用一个外部HTTP接口没有设置连接超时和读取超时用的是JDK默认的无限等待。下游某台机器网络抖动所有请求都卡在读响应上连接池被占满新的请求全部排队。消费线程拉取到消息后把任务丢给业务线程池线程池的队列越堆越长poll循环还在继续但消息实际上没人处理。这种问题在代码review时非常难发现因为代码逻辑没有bug问题出在非功能配置上。所以我在排查消费卡顿时会把这些点全部过一遍所有外部调用是否设置了connectTimeout和readTimeout默认值是什么连接池大小是否够用队列有没有界拒绝策略是什么异常处理里有没有吞异常——catch之后只有一行log.info甚至什么都不写消息处理里有没有隐式的大事务、全局锁、慢SQL消费线程池和业务线程池是不是同一个有没有相互拖累。根因找到之后修复方案也要分三层第一层是给所有外部调用设合理超时并加上失败降级第二层是把消费线程池与业务线程池隔离避免一个慢任务拖垮所有消费第三层是加告警当线程池队列深度超过阈值、消费速率低于预期时自动报警。没有第三层这个Bug迟早换一种方式再来一次。3. 横向迁移那些长得不一样的Bug底层原因却惊人的相似3.1 Pulsar Key_Shared模式不消费从分区到Key的思维转换有一次同事来问我Pulsar的Key_Shared模式突然不消费了消费者都连着但消息就是卡着不动。有了Kafka的排查经验我第一反应是看每个消费者是否都在干活。但Pulsar Key_Shared与Kafka分区模式有一个关键差异在Key_Shared模式下同一条Key的所有消息会固定发送到同一个消费者而且按顺序处理。这意味着只要某个Key的消息处理被阻塞这条Key上的所有后续消息都会堵住即使其他消费者是空闲的。排查方向完全不一样了。不要只看整体消费速率要先看每个消费者上是否堆积了大量pending消息有没有某个Consumer的队列被打满。再把积压消息的Key取出来去代码里搜这个Key对应的处理逻辑看看是不是遇到了慢接口、死循环或者幂等外键冲突。你会发现它和Kafka案例的底层逻辑是同一个某个消费者线程被特定业务卡死导致整个队列停摆。Pulsar的客户端配置也值得检查比如receiverQueueSize如果太小消费吞吐受限如果利用key_shared实现严格顺序务必确保生产者设置正确的Key否则会均匀散落到不同消费者上顺序语义被打破排查难度直接翻倍。3.2 Java容器内存飙高先分清堆内还是堆外容器里跑Java服务内存持续上涨直到被OOMKilled这是云原生时代最常见的Bug之一。很多人一上来就jmap -dump堆内存结果发现堆占用并不高问题出在堆外。所以第一条经验是不要凭直觉选工具先判断是堆内还是堆外。堆内问题用jmap -heap看堆使用率、jmap -histo看对象分布或者直接生成Heap Dump后用MAT分析。常见原因是缓存无上限、线程池队列无界、一次性加载了大集合。堆外问题则要复杂一些可能包括DirectByteBuffer未释放、JNI本地内存泄漏、元空间无限增长。JDK 8以上可以开启Native Memory Tracking-XX:NativeMemoryTrackingsummary运行中执行jcmd pid VM.native_memory summary能看到堆内、元空间、编译、GC、线程、代码缓存等各分区的内存变化。如果线程数量爆炸每个线程默认栈大小1MB几百个线程就是几百MB这种问题在线程Dump里一眼就能看出来。还有一种很隐蔽的情况容器内存Limit限制的是Cgroup但JVM默认不感知Cgroup导致JVM按宿主机内存分配堆大小。Kubernetes环境下如果没有正确配置UseContainerSupport会带来一大堆内存误判问题。这些问题表面上是内存高但排查的底层思路仍然是那一套先看数据分桶再逐层定位。3.3 硬件与网络里的隐蔽坑IP冲突、CAN终端电阻与核显Reset不是只有后端服务才会遇到时好时坏的Bug硬件和网络领域的排查思路很多可以复用但细节完全不同。IP冲突是我排查过最玄学的一种。现象是网络时通时断ping有时候通有时候不通远程连接频繁断开。用arpwatch和抓包软件看ARP请求发现局域网里有两个设备在抢同一个IP它们交替响应导致数据包被路由到错误的设备。排查步骤很简单先ping目标IP确认通断规律再查交换机ARP表用arping确认MAC地址是否变化。一旦确定IP冲突静态绑定IP和MAC或者给一台设备换IP就解决。CAN通信的故障排查里终端电阻是最容易被忽略的。CAN总线标准要求在物理总线两端各接一个120Ω终端电阻用来匹配阻抗、防止信号反射。如果两个终端电阻都缺失总线上的电平会异常表现为距离近的时候通信正常线拉长或者节点多了就开始丢帧。用示波器看CAN_H和CAN_L的差分波形正常情况下应该是0到2V左右的清晰方波缺了终端电阻波形边缘会出现明显的振铃。这不是软件Bug但排查思路和软件一样先确认物理层再看协议层逐层排除。AMD核显Reset Bug则是另一个方向的例子。某些平台的核显在休眠唤醒或高负载时会发生reset表现是屏幕黑一下、驱动进程崩溃日志里会出现GPU reset相关报错。排查时先通过dmesg看reset时间点再尝试更新内核、关闭某个电源管理特性、调整BIOS里的显存分配。这类问题往往跟特定驱动版本强相关换一台机器可能永远复现不了——和软件里换个环境就好的偶发Bug非常像。3.4 前端Bug排查从Console到断点调试的实战姿势前端的Bug和上面那些不一样很多不是逻辑错而是你以为是A问题实际是B问题。我见过不少前端同事排查问题只会console.log打了几十个log还找不到原因最后发现是数据结构不对。前端调试的正确姿势应该是分层推进。第一步是看Network面板确认请求有没有发出、响应码是什么、耗时多少、有没有被缓存。很多时候前端Bug实际上是接口报错或者跨域配置问题这一步能排除掉一半的假Bug。第二步是根据报错信息直接在Sources面板打断点而不是到处加console.log。打断点有个技巧在怀疑函数入口和关键分支都打上然后单步执行观察作用域里的变量。比log高效得多。第三步是善用条件断点和Watch表达式。比如在一个循环里只对某个特定值停住条件断点可以极大节省时间。对于线上偶发问题可以配合Performance面板录制一段操作看主线程和任务执行时间线很多卡顿和白屏问题从火焰图上一眼就能看出来。说实话前端和后端排错的底层逻辑完全一致收集现场、分层定位、验证假设。4. 最让人头疼的偶发Bug无法复现时怎么办4.1 无法复现的本质变量太多我在本地跑是好的 测试环境也复现不了 上了线偶尔才出现一次——这类话出现的时候就是最难熬的时刻。无法复现不是Bug不存在而是触发条件太隐蔽。我把它归为三类环境差异、并发时序、脏数据。环境差异包括JDK版本不同、依赖库版本不同、系统时区不同、容器资源限制不同。比如日期格式化在JDK8和JDK11解析行为有差异测试环境是Java11线上是Java8偶发报错当然复现不了。并发时序是最常见的元凶典型的HashMap在多线程扩容时形成循环链表导致某个线程CPU飙到100%。这种Bug反复压测不一定能触发因为需要精确的扩容时机。脏数据则是指数据库里存在某些特殊值——某个字段Null、某个字符串超长、某个时间戳是1970年——代码在写的时候没考虑这些边界生产环境一遇到就出问题。面对这类问题我的原则是先用最小代价扩大信息量而不是急着找复现路径。只要你把现场描述清楚把异常参数记录下来即使复现不了也能通过代码审查和压力测试缩小怀疑范围。4.2 用增强日志和埋点让Bug自证其罪既然不能靠运气复现那就让系统在下次发生时自动留下证据。最有效的手段是增强日志。在生产环境的可疑路径上加上带traceId、业务参数、耗时、重试次数的结构化日志。注意不是所有地方都加要加在你能想象到的如果这里出问题一定需要什么信息的地方。如果不想发版可以用动态日志级别调整。Arthas的logger命令可以直接在运行中的服务上修改某个类的日志级别把debug日志临时打开等现场信息收集完再关掉。类似的方式还有log4j2的运行时修改。更推荐的做法是在关键业务节点加Metrics计数器比如记录消息处理次数平均耗时异常次数线程池队列深度。当偶发问题再次出现时这些数据能让你快速定位到具体阶段。我一直认为好的可观测性建设不是给监控平台看的是为了让Bug在发生后能自证其罪。4.3 保住现场线程Dump、Heap Dump与流量回放偶发Bug最大的敌人是重启。服务一重启现场就没了所有线索归零。所以遇到无法复现的问题第一要务是尽可能把现场保留下来而不是急着恢复。核心转储、线程Dump、Heap Dump、日志快照能拿多少拿多少。这里分享一个我个人的做法给关键服务加一个遇险自拍机制。通过在代码里埋一个健康探针当线程池队列达到阈值、CPU持续飙高、接口耗时超过N秒时自动执行jstack、jmap -dump或者记录关键线程的堆栈。虽然要写一点代码但它是应对偶发Bug最稳定的保底手段。流量回放也很有用。比如用GoReplay或Tcpcopy录制生产环境的请求流量在测试环境重放很多并发时序问题能被压出来。另一种方式是构造一个和生产数据特征一致的数据集专门测试边界条件空值、超长字符串、并发高峰、慢网络。当你在这些条件下复现出Bug时往往也就找到了根因。5. 把一次排查变成团队的排查能力5.1 一份合格的Bug排查报告应该长什么样排查完成之后写报告这个动作不能省。但很多人写的报告只写了结论——某接口超时导致消费阻塞已修复。这对团队几乎没有复用价值。我更看重过程一份合格的报告至少要有五个部分背景与现象什么时间、什么告警、影响范围多大排查时间线每一个步骤做了什么、得出了什么判断、排除了什么可能根因分析为什么会出现这个问题是技术债、配置缺失还是设计缺陷修复方案与验证改了什么、怎么验证的、有没有回归风险预防措施加什么监控、补什么checklist、代码评审要关注什么。技术上根因可能很简单但过程才是团队能复用和演练的东西。我见过有经验的工程师能把一个半小时的排错过程浓缩成一页纸的流程图后来的人照着这个思路十来分钟就能定位同类问题。5.2 建立Checklist与Runbook让经验不依赖个人如果团队里排Bug的能力集中在一两个人身上是极其危险的。把常见问题整理成Runbook和Checklist能让整个团队在面对告警时不慌、不漏步骤。我建议按问题域建立几个清单维护在团队知识库里消息积压排查清单Lag上涨→看消费速率→看消费者组→看线程Dump→看外部调用超时内存异常排查清单先分堆内堆外→堆内看Dump→堆外看NMT→确认容器内存配置网络通信故障清单先看ARP/IP冲突→抓包确认二层通信→再排查三层路由硬件偶发故障清单先记录日志时间点→对比驱动版本→尝试复现条件Linux主机异常排查清单先看登录日志和认证日志再查可疑进程、定时任务、启动项和文件变更。这些清单不需要多完美哪怕一开始只有粗略框架在一次次实战中不断添加新坑半年后就会成为团队最值钱的经验库。核心思想是把排错从一个依赖个人经验的手艺活变成一个团队通用的标准流程。5.3 我在多次排查之后养成的几个习惯最后聊一点个人习惯。我做Bug排查这些年最有价值的改变是从着急修变成先取证。发现Bug后第一件事不再是定位和修改而是先把现象、时间、影响面、日志、指标存下来。这就像事故现场的拍照取证证据完整了后面不管怎么折腾都不会慌。第二个习惯是写排查日记。不是写技术总结而是记录当时的思路和情绪——我起初以为是内存泄漏花了二十分钟看Heap Dump后来发现是GC配置问题。半年后再翻能清楚看到自己的思维盲区。这种错题本比任何培训都管用。第三个习惯是先怀疑配置再怀疑代码。线上很多Bug不是代码写错而是配置没有跟上超时没设置、线程池太小、日志格式不对、JVM参数不合理。代码评审时多关注非功能配置比抠某一行逻辑更能避免线上事故。
返回列表