Kafka消息消费与线程池堵塞排查优化
Kafka消息消费与线程池堵塞排查优化
生产环境中出现 Kafka 消息积压时,第一反应往往是“Kafka 堵了”。但多数情况下,真正的问题不在 Kafka,而在消费者处理速度跟不上生产速度。
以告警消息为例,一条消息里可能包含多个测点,每个测点又可能涉及设备信息获取、告警规则匹配、告警记录落库、Redis 状态刷新、通知推送等动作。只要其中某一步变慢,Kafka 消费线程就会被拖住,最终表现为 consumer group lag 增长。
本文参考项目中 listenAlarm 的处理逻辑,不展开具体业务代码,重点说明线程池排查思路、Kafka 堵塞排查方法和优化方向。
什么时候需要排查
这类问题最明显的表象通常不是服务直接报错,而是时效性变差。比如消息已经产生了,但业务侧很久之后才看到处理结果。
常见表象包括:
| 表象 | 说明 | 优先怀疑方向 |
|---|---|---|
| 告警时效性变差 | 设备已经触发告警,但页面、通知或记录明显延迟 | Kafka 消费慢、主链路处理慢 |
| Kafka lag 持续增长 | consumer group 的 LAG 越来越大 | 消费速度低于生产速度 |
| 告警记录入库延迟 | 消息到了,但数据库中的告警记录生成慢 | 数据库写入慢、规则匹配慢 |
| 页面状态刷新慢 | 告警记录已有,但页面状态迟迟不变 | Redis 刷新慢、异步线程池堵塞 |
| 通知推送延迟 | 告警已生成,但短信、站内信、浏览器通知很久才到 | 通知线程池堵塞、外部通知接口慢 |
| 服务 CPU 升高 | 消费服务 CPU 长时间偏高 | 规则计算、日志、序列化、循环处理 |
| 服务内存上涨 | 堆内存持续增长,GC 频繁 | 无界队列堆积、消息对象积压 |
| 日志出现超时 | 远程调用、Redis、数据库偶发超时 | 下游资源慢或连接池不足 |
只要出现“消息处理结果明显滞后”这类时效性问题,就应该按 Kafka 积压、消费线程、线程池、下游资源这几个方向逐层排查。
参考的业务处理链路
listenAlarm 的整体逻辑可以抽象成下面这条链路:
Kafka 消费线程收到告警消息
-> 解析消息体
-> 遍历消息中的多个测点
-> 同步处理每个测点的告警逻辑
-> 查询或注册设备信息
-> 获取设备对应的告警规则
-> 创建或维护相关表结构、标签信息
-> 匹配告警规则
-> 写入告警记录
-> 异步刷新 Redis 状态
-> 异步推送告警通知
这里最关键的一点是:主告警处理逻辑是在 Kafka 消费线程里同步执行的。也就是说,当前消息里的测点没有处理完,listener 方法就不会结束,当前分区的 offset 推进也会变慢。
后面的 Redis 状态刷新、通知推送虽然是异步提交到线程池中执行,但主链路里的设备处理、规则匹配、告警落库仍然会影响 Kafka 消费速度。
为什么单条消息慢会导致 Kafka 积压
假设一条 Kafka 消息中有 20 个测点,每个测点平均处理 2 秒,那么这条消息完整处理可能需要 40 秒。这个过程中 Kafka 消费线程一直在执行业务逻辑,无法快速拉取并处理后续消息。
堵塞过程通常是:
单个测点处理慢
-> 一条 Kafka 消息整体处理时间变长
-> Kafka listener 线程长时间被业务占用
-> 当前分区 offset 推进变慢
-> consumer group lag 持续增长
-> 告警消息延迟越来越高
如果业务里还有异步线程池,比如告警状态刷新、通知推送线程池,那么还可能出现第二层问题:Kafka 主消费链路已经变慢,同时异步线程池队列也在堆积,导致 JVM 内存上涨、任务延迟变大。
需要重点怀疑的耗时点
排查时不要只盯着线程池参数,要先拆解一条消息到底慢在哪里。
| 环节 | 可能问题 | 现象 |
|---|---|---|
| 消息解析 | 消息体过大,测点数量多 | listener 单次执行时间变长 |
| 详细日志 | 打印完整测点内容、日志量过大 | CPU 和磁盘 IO 升高 |
| 设备注册或查询 | 远程服务慢、超时过长、连接池不足 | 线程栈停在 HTTP/RPC 调用 |
| 告警规则获取 | 缓存未命中、Redis 慢、DB 查询慢 | 单测点处理耗时不稳定 |
| 表结构和标签处理 | 建表、改标签、初始化逻辑频繁执行 | TDengine 或数据库压力升高 |
| 告警规则匹配 | 规则数量多、表达式复杂 | CPU 升高或处理时间变长 |
| 告警记录写入 | 批量写入慢、连接不足、下游抖动 | 线程栈停在 DAO 或 JDBC |
| Redis 状态刷新 | Redis 慢、连接池满 | 异步线程池队列增长 |
| 通知推送 | 浏览器通知、短信、外部接口慢 | 通知线程长时间执行 |
判断慢点的核心方法是:加分段耗时日志,或者在线上用 Arthas、jstack 看线程长期停在哪里。
线程池相关风险
在告警处理链路中,常见的线程池风险主要有三类。
第一类是消费线程同步处理太多业务逻辑。Kafka listener 线程本身就是有限资源,如果在里面做大量 DB、Redis、HTTP 操作,任何一个下游变慢都会直接拖慢 Kafka 消费。
第二类是异步线程池队列无界。如果 Redis 刷新、通知推送等任务提交到无界队列,任务生产速度大于消费速度时,队列会一直增长。短时间看没有报错,但延迟和内存都会持续上升。
第三类是不同任务共用一个线程池。实时告警、状态刷新、通知推送、历史修复任务如果混用同一个线程池,某一类慢任务会影响其他任务。
优化方向一:先补齐耗时日志
优化前先定位。建议至少记录这些耗时:
| 日志点 | 目的 |
|---|---|
| 单条 Kafka 消息总耗时 | 判断 listener 是否被长期占用 |
| 每个测点处理耗时 | 找出是否存在个别慢测点 |
| 设备注册或查询耗时 | 判断远程服务是否拖慢主链路 |
| 告警规则获取耗时 | 判断规则缓存或查询是否异常 |
| 表结构和标签处理耗时 | 判断建表、标签维护是否频繁且慢 |
| 告警规则匹配耗时 | 判断规则表达式是否过多或过复杂 |
| 告警记录落库耗时 | 判断数据库或 TDengine 是否为瓶颈 |
| Redis 状态刷新耗时 | 判断异步线程池是否被 Redis 拖慢 |
| 通知推送耗时 | 判断通知接口是否拖慢线程池 |
日志中建议带上 recordKey、subSn、objId、测点数量、线程名和耗时。不要打印过大的完整消息体,否则排查日志本身也会成为性能问题。
优化方向二:减少 Kafka 消费线程里的慢操作
Kafka listener 中应尽量只做必要处理,避免长时间等待外部资源。
可以从这些方向优化:
- 减少消费线程里的大对象转换和完整报文日志。
- 对远程调用设置明确超时时间,避免线程长时间卡住。
- 设备信息、告警规则优先走缓存,减少每条消息都查库或远程调用。
- 表初始化、标签初始化尽量前置,不要放在高频消费路径反复执行。
- 告警记录尽量批量写入,减少频繁单条写库。
- 无效消息、心跳消息、无告警规则的消息尽早返回。
- 主链路只做核心告警判定和落库,非核心通知类动作异步处理。
如果单条消息里测点很多,还要关注批次大小。一次消息处理太多测点,会让单次 listener 执行时间过长。
优化方向三:谨慎把主链路异步化
一种常见优化是:Kafka listener 只负责解析消息,然后把每个测点提交到业务线程池处理。这样 listener 可以更快返回,Kafka lag 可能会下降。
但这个方案有一个重要风险:listener 返回后,Kafka offset 可能已经提交,而线程池里的业务任务还没有真正处理完成。如果服务重启或线程池任务失败,就可能出现“Kafka 认为已消费,业务实际没处理完”的问题。
所以主链路异步化前要先设计好:
- 是否手动提交 offset。
- 是否等业务处理完成后再确认消费成功。
- 业务是否具备幂等处理能力。
- 失败任务是否有补偿机制。
- 线程池满了以后是暂停消费、拒绝并告警,还是写入补偿队列。
- 告警消息是否允许丢弃,通常告警消息不建议静默丢弃。
异步化不是简单加一个线程池,而是消费语义、失败重试、幂等和补偿一起设计。
优化方向四:线程池必须有界并可监控
异步线程池建议使用有界队列,避免任务无限堆积。
需要重点监控这些指标:
| 指标 | 异常表现 | 含义 |
|---|---|---|
| activeCount | 长期接近最大线程数 | 工作线程一直满负荷 |
| queueSize | 持续增长 | 提交速度大于处理速度 |
| completedTaskCount | 增长很慢 | 单个任务耗时长或卡住 |
| rejectedCount | 出现拒绝 | 系统超过设计容量 |
| taskCost | 单任务耗时升高 | 下游或业务逻辑变慢 |
线程池拒绝策略不要静默丢弃。至少要记录日志、打点监控、触发告警,并结合业务决定是否重试或补偿。
优化方向五:拆分不同类型任务
不要把所有告警相关任务都放进同一个线程池。更合理的拆分方式是:
告警主处理线程池 处理核心告警判定和落库
Redis状态刷新线程池 处理设备告警状态缓存刷新
通知推送线程池 处理浏览器通知、短信、站内信
历史修复线程池 处理历史数据刷新、批量修复任务
这样做的好处是:通知接口慢不会拖垮 Redis 状态刷新,历史任务也不会抢占实时告警资源。
线程池拆分后,每个线程池都要有独立的队列大小、拒绝策略、耗时日志和监控指标。
优化方向六:调整 Kafka 并发和分区
如果主链路已经优化过,但消费能力仍然不足,可以考虑调整 Kafka 并发。
需要注意:同一个 consumer group 中,一个分区同一时间只能被一个消费者线程消费。如果 topic 只有 1 个分区,单纯把 listener concurrency 调大并不会明显提升消费能力。
调整顺序建议是:
- 先看 topic 分区数。
- 再看 consumer 实例数和 listener 并发数。
- 确认消息 key 是否导致数据集中到少数分区。
- 根据机器资源和下游承载能力提高并发。
- 压测验证,避免把数据库、Redis、远程服务打满。
并发不是越大越好。如果下游已经慢,继续提高消费并发只会把下游压得更慢。
生产排查第一步:看 Kafka 是否积压
先查 consumer group lag:
kafka-consumer-groups.sh \
--bootstrap-server 127.0.0.1:9092 \
--describe \
--group xxx-consumer-group
重点看这些字段:
| 字段 | 说明 |
|---|---|
CURRENT-OFFSET | 当前消费到的位置 |
LOG-END-OFFSET | 分区最新位置 |
LAG | 积压消息数 |
CONSUMER-ID | 消费者实例 |
HOST | 消费者所在机器 |
CLIENT-ID | 客户端标识 |
如果目标 topic 的 LAG 持续增长,说明消费速度跟不上生产速度。
再看 topic 分区:
kafka-topics.sh \
--bootstrap-server 127.0.0.1:9092 \
--describe \
--topic xxx-topic
如果分区数很少,增加消费者实例或调大并发效果有限。
生产排查第二步:定位 Java 进程
jps -l
或者:
ps -ef | grep xxx-service
拿到 PID 后,看整体 CPU 和内存:
top -p <pid>
如果 CPU 不高但 Kafka lag 增长,通常是线程在等待 DB、Redis、HTTP、锁或 IO。
如果 CPU 很高,继续看哪个线程占 CPU。
生产排查第三步:看哪个线程 CPU 高
查看 Java 进程内线程 CPU:
top -H -p <pid>
找到 CPU 高的线程 ID 后,转成 16 进制:
printf '%x\n' <tid>
导出线程栈:
jstack -l <pid> > xxx.jstack
按 16 进制线程 ID 查:
grep -n "nid=0x线程ID十六进制" xxx.jstack -A 80
如果高 CPU 线程是 Kafka listener 线程,说明消费线程本身在忙。若是告警异步线程池线程,说明异步任务处理慢或任务量过大。
生产排查第四步:看线程是否长时间卡住
连续抓 3 次线程栈,每次间隔 10 秒:
jstack -l <pid> > jstack-1.txt
sleep 10
jstack -l <pid> > jstack-2.txt
sleep 10
jstack -l <pid> > jstack-3.txt
搜索关键业务方法或线程名前缀:
grep -n "listenAlarm" jstack-1.txt -A 80
grep -n "deviceAlarm" jstack-1.txt -A 80
grep -n "告警" jstack-1.txt -A 80
grep -n "Redis" jstack-1.txt -A 80
如果 3 次线程栈里同一个线程一直停在相同位置,说明该位置耗时异常。
常见判断:
| 栈位置 | 可能原因 |
|---|---|
| HTTP/RPC 调用 | 远程服务慢或超时设置过长 |
| Redis 调用 | Redis 慢、连接池满、网络抖动 |
| DAO/JDBC/TDengine 调用 | 查询慢、写入慢、连接不足 |
| BLOCKED 状态 | 锁竞争严重 |
| WAITING/TIMED_WAITING | 等待资源、连接池、队列或锁 |
| RUNNABLE 且 CPU 高 | 计算、序列化、日志、规则匹配等 CPU 消耗高 |
生产排查第五步:使用 Arthas 看方法耗时
启动 Arthas:
java -jar arthas-boot.jar
看最忙线程:
thread -n 20
观察 Kafka 消费方法耗时:
trace xxx.xxx.XxxListener xxxListenMethod '#cost > 1000'
观察主处理方法耗时:
trace xxx.xxx.XxxListener xxxHandleMethod '#cost > 1000'
观察设备告警处理耗时:
trace xxx.xxx.XxxServiceImpl xxxBusinessMethod '#cost > 1000'
观察参数、异常和耗时:
watch xxx.xxx.XxxServiceImpl xxxTargetMethod '{params, returnObj, throwExp, #cost}' '#cost > 1000' -x 2
Arthas 在生产上要控制范围,优先加 #cost > 1000 这类条件,避免输出过多。
生产排查第六步:看下游资源是否拖慢
线程卡住通常不是线程本身的问题,而是线程等待下游资源。
可以检查连接和 GC:
# 查看进程连接数
netstat -antp | grep <pid> | wc -l
# 查看 Redis 连接
netstat -antp | grep 6379
# 查看数据库或 TDengine 连接,端口按实际环境替换
netstat -antp | grep 6030
# 查看打开文件数
ls /proc/<pid>/fd | wc -l
# 查看 GC 情况
jstat -gcutil <pid> 1000 10
# 查看堆信息
jcmd <pid> GC.heap_info
如果 GC 频繁、堆持续上涨,同时线程池队列无界或队列持续增长,重点怀疑异步任务堆积。
线上处理顺序
遇到 Kafka 告警消息积压时,建议按这个顺序排查:
- 用
kafka-consumer-groups.sh确认目标 topic 是否 lag 持续增长。 - 用
top -H -p <pid>找高 CPU 或长时间运行线程。 - 用
jstack判断线程卡在消费方法、主业务处理、Redis、数据库还是通知推送。 - 用 Arthas
trace定位最慢的方法。 - 先优化最慢的业务调用,再考虑增加 Kafka 并发或线程池大小。
- 检查异步线程池是否队列堆积,必要时拆分线程池。
- 检查下游数据库、Redis、远程服务是否已经达到瓶颈。
最推荐的改造方向
结合这类告警消费链路,优先级建议如下:
- 第一优先级:补齐分段耗时日志,确认慢点在哪里。
- 第二优先级:给远程调用、Redis、数据库操作设置合理超时。
- 第三优先级:异步线程池改成有界队列,并增加队列长度、活跃线程数、完成任务数监控。
- 第四优先级:减少 Kafka listener 线程里的慢操作,把可延后的通知、状态刷新异步化。
- 第五优先级:实时任务、通知任务、历史任务拆分线程池。
- 第六优先级:如果主处理也要异步化,必须配套手动 offset、幂等和补偿机制。
- 第七优先级:根据 topic 分区数和机器资源,调整 Kafka 并发和消费者实例数。
小结
Kafka 消息积压时,不能只看 Kafka,也不能只调线程池参数。正确排查顺序是:先确认 lag,再定位消费线程是否被业务占用,然后看异步线程池是否堆积,最后检查数据库、Redis、远程服务等下游资源。
线程池优化的重点不是简单加线程,而是让任务有边界、有监控、有超时、有隔离、有补偿。这样生产环境出现堵塞时,才能快速判断是 Kafka 消费慢、线程池堵,还是下游服务拖慢。
