生产环境Kafka消息疯狂积压,我排查了一整夜,总结出这套保命流程
消息积压到底意味着什么
先说个很多人容易搞混的事情。Kafka消息积压,表面上看是消费者消费不过来了,但实际情况远比这个复杂。
积压的本质就是:生产速度 > 消费速度,或者消费者压根就不工作了。
听起来是废话对吧?但你仔细想想,导致这个结果的原因可能有十几种。生产端突然流量暴增、消费端代码有Bug卡住了、消费者实例挂了几个、下游数据库慢了拖住了消费线程、甚至Kafka集群自身有Broker出了问题……每一种情况的处理方式都不一样。
我见过最离谱的一次,积压原因居然是有个同事在消费逻辑里加了一个Thread.sleep(1000)用来调试,然后忘了删就发上去了。一个sleep搞崩整条链路,你说气不气人。
所以遇到积压,千万别上来就想着加消费者实例或者重启,先搞清楚到底是哪个环节出了问题。
出事了先看什么
那天晚上我其实也慌了一小会,但后来慢慢形成了一套固定的排查顺序,每次都照着来,效率高很多。
第一件事:确认积压的范围。
用kafka自带的命令看一下Consumer Group的消费情况:
kafka-consumer-groups.sh --bootstrap-server kafka01:9092 --describe --group your-consumer-group这个命令会把每个Partition的当前Offset、Log End Offset、Lag都列出来。重点看Lag这一列,哪些Partition积压严重一目了然。
有时候你会发现,不是所有Partition都积压,可能就那么两三个Partition的Lag特别高。这就说明问题大概率出在这几个Partition对应的消费者实例上。
然后看消费者实例还活着没。
同样用上面那个命令,看CONSUMER-ID和HOST列。如果某些Partition没有分配到消费者,说明消费者实例挂了或者发生了Rebalance。
我遇到过好几次,消费者因为处理时间太长触发了session.timeout导致被踢出Group,然后不断Rebalance,消费基本就停了。日志里一直刷Rebalance相关的信息,消费吞吐量直接归零。
# 看消费者日志里有没有这种关键字
grep -i "rebalance\|kicked\|timeout\|leaving group" /path/to/consumer.log接着看Kafka集群本身。
有时候不是消费者的锅。Broker的磁盘IO打满了,或者某个Broker直接宕了导致Leader切换,ISR缩减,这些都会间接影响消费速度。
# 查看Topic的Partition分布和ISR情况
kafka-topics.sh --bootstrap-server kafka01:9092 --describe --topic your-topic重点看Isr列表是不是和Replicas一致。如果ISR比Replicas少,说明有Broker掉队了。这种情况下Broker日志得仔细看。
最容易误判的几个坑
说几个我自己和身边同事踩过的坑,真的是血泪教训。
坑一:只看Consumer Lag数字,不看变化趋势。
Consumer Lag是60万,严重不严重?不一定。如果Lag在持续下降,说明消费者在追数据,可能很快就追上了,这时候你去重启反而帮倒忙。但如果Lag在持续上升或者纹丝不动,那才是真出事了。
所以一定要看Lag的变化曲线,不要只看某一个时间点的绝对值。我们后来在Grafana上专门做了一个Consumer Lag Rate的面板,看Lag的增减速率,比看绝对值有用太多了。
坑二:消费者假死。
消费者进程还在,端口还通,但实际上内部线程已经卡死了。这种情况下你用kafka-consumer-groups.sh看,会发现这个消费者还挂在Group里,但它的Offset就是不动。
我有一次排查了快一个小时才发现是这个问题。消费者里调了一个外部HTTP接口,那个接口没有设超时时间,正好那天对方服务挂了,线程全卡在HTTP请求上。连接池耗尽,后面的消息全部阻塞。
所以消费者里调外部服务,一定一定要设超时时间,这个真的是铁律。
// 反面教材:没有超时
HttpResponse response = httpClient.execute(request);
// 正确做法:设置连接超时和读取超时
RequestConfig config = RequestConfig.custom()
.setConnectTimeout(3000)
.setSocketTimeout(5000)
.build();
HttpGet request = new HttpGet(url);
request.setConfig(config);坑三:Partition数量和消费者实例数不匹配。
这个说出来好像很基础,但真的有人踩。Kafka里一个Partition同一时刻只能被同一个Consumer Group里的一个消费者消费。如果你Topic有10个Partition,但只起了3个消费者实例,那剩下的7个Partition就要分摊到这3个实例上。消费能力严重不足。
反过来,你起了20个消费者,但Partition只有10个,那多出来的10个实例就是干等着,完全空闲。
这个关系一定要搞清楚:消费者实例数 ≤ Partition数才有意义,最理想的情况是1:1。
一套标准的排查流程
吃了几次亏之后,我把排查流程整理成了一个Checklist,贴在工位上,后来也分享给了组里的同事。这里也分享出来:
Step 1:确认积压范围和程度
- 用kafka-consumer-groups.sh看Lag分布
- 区分是全量积压还是个别Partition积压
- 看Lag是在涨还是在降
Step 2:检查消费者状态
- 消费者进程是否存活
- 是否在频繁Rebalance
- JVM是否有Full GC(Java消费者)
- 消费线程是否有阻塞/死锁
# Java应用查看线程状态
jstack <pid> | grep -A 20 "kafka-consumer"
# 查看GC情况
jstat -gcutil <pid> 1000 10Step 3:检查消费逻辑
- 有没有慢SQL拖住消费
- 有没有调外部接口没设超时
- 有没有异常处理不当导致重试风暴
- 批量消费的batch size是否合理
Step 4:检查Kafka集群
- Broker是否有宕机
- ISR是否正常
- 磁盘IO和网络IO是否打满
- Controller是否正常
# 查看集群中的Controller
kafka-metadata.sh --snapshot /path/to/metadata --controller
# 或者通过ZK查看(老版本)
echo dump | nc localhost 2181 | grep controllerStep 5:检查生产端
- 生产流量是否异常暴增
- 有没有消息倾斜(某几个Key的消息集中到少数Partition)
Step 6:临时处理 + 根因修复
- 临时:扩消费者实例、跳过异常消息、降低处理逻辑复杂度
- 根因:修复代码Bug、优化慢查询、调整超时配置、扩Partition
那次P0故障的完整复盘
回到文章开头那次事故,我把完整过程讲一下。
当晚11点15分左右收到告警,我11点半左右上线开始排查。
先用kafka-consumer-groups.sh看了一下,发现是订单相关的Topic积压最严重,Lag已经到了40多万而且还在涨。其他几个Topic基本正常。
然后看消费者状态,发现5个消费者实例都还活着,没有Rebalance的迹象。但是Offset几乎不动,每分钟才前进几十条,正常情况下应该是每秒几千条。
这就奇怪了,实例都在但不消费。我上机器看了一下日志,发现大量的超时报错,连接的是下游的一个订单状态同步接口。
然后去问了下游团队,才知道他们那边在做数据库迁移,接口响应时间从原来的几十毫秒飙到了好几秒,甚至有的请求直接超时了。
我们的消费逻辑是同步调用这个接口的,一条消息消费要等接口返回才能处理下一条。接口慢了,消费自然就堵住了。
临时处理:
先把那个接口调用的超时时间从10秒降到3秒,快速失败。然后把失败的消息写到一个重试Topic里,不阻塞主流程。同时把消费者实例从5个扩到10个(Topic有20个Partition,之前5个其实就不太够)。
这几步下来,大概12点半的时候Lag开始明显下降了,到凌晨2点基本追平。
事后整改:
- 消费逻辑里所有外部调用统一加超时,最长不超过3秒
- 引入异步消费模式,消费线程和业务处理线程分离
- 加了一个死信队列(DLQ),处理失败超过3次的消息扔进去,人工处理
- 下游接口加了熔断,用的Resilience4j,接口失败率超过50%直接熔断,不再请求
- 监控补齐:Consumer Lag的告警阈值从10万降到1万,告警分级
最后写了个复盘文档,在周会上讲了一遍。说实话那天晚上确实挺狼狈的,但经过这次之后整个消费链路稳定性提升了不少。
几个避坑建议
根据我这两年的经验,关于Kafka消费积压,有几条铁律分享一下:
不要在消费逻辑里做重计算。 消费者应该尽量轻,拿到消息之后快速分发到业务线程池去处理。消费线程本身不要做太重的事情。
不要忽略max.poll.interval.ms这个参数。 默认值是5分钟,意思是如果两次poll之间超过5分钟,Kafka就认为你这个消费者挂了,会触发Rebalance。如果你的单批处理时间比较长,记得把这个值调大。
# 消费者关键参数参考
max.poll.records=500
max.poll.interval.ms=600000
session.timeout.ms=30000
heartbeat.interval.ms=10000
fetch.min.bytes=1
fetch.max.wait.ms=500不要一出事就想着重置Offset。 我见过有人一遇到积压就把Offset重置到Latest,消息直接跳过了。短期来看Lag确实没了,但跳过的那些消息里可能有重要的业务数据,后面对账的时候才发现少了一堆,那才是真正的灾难。
不要省略监控和告警。 Consumer Lag、消费TPS、消费延迟、Rebalance次数,这几个指标一定要有。别等到用户投诉了才发现积压了几个小时。
下面这个Prometheus+JMX的监控配置可以参考:
# kafka consumer JMX metrics
- pattern: "kafka.consumer<type=consumer-fetch-manager-metrics, client-id=(.+)><>records-lag-max"
name: kafka_consumer_records_lag_max
type: GAUGE
labels:
client_id: "$1"
- pattern: "kafka.consumer<type=consumer-fetch-manager-metrics, client-id=(.+)><>records-consumed-rate"
name: kafka_consumer_records_consumed_rate
type: GAUGE
labels:
client_id: "$1"不要在没有预案的情况下上线消费者变更。 改了消费逻辑、调了并发数、换了反序列化方式,任何变更都要有回滚方案。我们现在的规矩是消费者变更必须灰度,先上一个实例跑几个小时没问题再全量。
一些实用的小工具和脚本
分享几个我日常用的东西,排查的时候能提高不少效率。
快速查看所有Consumer Group的Lag:
#!/bin/bash
# 批量查看所有consumer group的lag
for group in $(kafka-consumer-groups.sh --bootstrap-server kafka01:9092 --list); do
echo "========== $group =========="
kafka-consumer-groups.sh --bootstrap-server kafka01:9092 --describe --group $group 2>/dev/null | awk 'NR>1{sum+=$5}END{print "Total Lag: "sum}'
done消费速率监控脚本:
#!/bin/bash
# 每10秒采集一次Lag,计算消费速率
GROUP="your-consumer-group"
BOOTSTRAP="kafka01:9092"
while true; do
LAG=$(kafka-consumer-groups.sh --bootstrap-server $BOOTSTRAP --describe --group $GROUP 2>/dev/null | awk 'NR>1{sum+=$5}END{print sum}')
echo "$(date '+%H:%M:%S') Total Lag: $LAG"
sleep 10
done快速定位消息倾斜(Partition热点):
kafka-consumer-groups.sh --bootstrap-server kafka01:9092 --describe --group your-group | sort -t' ' -k5 -nr | head -10如果发现某几个Partition的Lag远高于其他,很可能是消息Key分布不均匀导致的数据倾斜。这种时候要去看生产端的Key设计是不是有问题。
写在最后
Kafka消息积压这个问题,说难也不难,说简单也不简单。关键在于你有没有一套成体系的排查思路,而不是出事了到处乱查碰运气。
我的体会是,排查问题本身可能只需要10分钟,但知道该查什么、从哪里查,可能需要踩很多坑才能积累出来。
这篇文章里的排查流程和避坑建议,都是我和我们团队在实际生产环境中一次次事故里总结出来的。希望能帮你在遇到类似问题的时候少走一些弯路,少熬一些夜。
如果你觉得有用,帮忙点个赞或者转发一下,让更多做运维的朋友看到。有什么问题也欢迎在评论区交流,大家一起进步。
关注公众号「耕云躬行录」,我会持续分享生产环境的运维实战经验、故障复盘、排查技巧。都是一线踩出来的坑,不讲虚的,只聊能用上的。
个人博客:躬行笔记,更多技术文章和工具资源,欢迎来逛。
下一期准备写Kafka集群扩容与Partition迁移的实战踩坑,感兴趣的可以先关注,更新不迷路。