Kafka消息消费与线程池堵塞排查优化

黑色的灵眸大约 15 分钟工作常见问题总结JavaKafka线程池性能排查告警处理

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 拖慢
通知推送耗时判断通知接口是否拖慢线程池

日志中建议带上 recordKeysubSnobjId、测点数量、线程名和耗时。不要打印过大的完整消息体,否则排查日志本身也会成为性能问题。

优化方向二:减少 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 调大并不会明显提升消费能力。

调整顺序建议是:

  1. 先看 topic 分区数。
  2. 再看 consumer 实例数和 listener 并发数。
  3. 确认消息 key 是否导致数据集中到少数分区。
  4. 根据机器资源和下游承载能力提高并发。
  5. 压测验证,避免把数据库、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 告警消息积压时,建议按这个顺序排查:

  1. kafka-consumer-groups.sh 确认目标 topic 是否 lag 持续增长。
  2. top -H -p <pid> 找高 CPU 或长时间运行线程。
  3. jstack 判断线程卡在消费方法、主业务处理、Redis、数据库还是通知推送。
  4. 用 Arthas trace 定位最慢的方法。
  5. 先优化最慢的业务调用,再考虑增加 Kafka 并发或线程池大小。
  6. 检查异步线程池是否队列堆积,必要时拆分线程池。
  7. 检查下游数据库、Redis、远程服务是否已经达到瓶颈。

最推荐的改造方向

结合这类告警消费链路,优先级建议如下:

  • 第一优先级:补齐分段耗时日志,确认慢点在哪里。
  • 第二优先级:给远程调用、Redis、数据库操作设置合理超时。
  • 第三优先级:异步线程池改成有界队列,并增加队列长度、活跃线程数、完成任务数监控。
  • 第四优先级:减少 Kafka listener 线程里的慢操作,把可延后的通知、状态刷新异步化。
  • 第五优先级:实时任务、通知任务、历史任务拆分线程池。
  • 第六优先级:如果主处理也要异步化,必须配套手动 offset、幂等和补偿机制。
  • 第七优先级:根据 topic 分区数和机器资源,调整 Kafka 并发和消费者实例数。

小结

Kafka 消息积压时,不能只看 Kafka,也不能只调线程池参数。正确排查顺序是:先确认 lag,再定位消费线程是否被业务占用,然后看异步线程池是否堆积,最后检查数据库、Redis、远程服务等下游资源。

线程池优化的重点不是简单加线程,而是让任务有边界、有监控、有超时、有隔离、有补偿。这样生产环境出现堵塞时,才能快速判断是 Kafka 消费慢、线程池堵,还是下游服务拖慢。

Loading...