Kafka 消息堆积,是生产环境里最容易遇到、也最容易“手忙脚乱”的问题之一。很多同学第一反应就是:堆了?加消费者啊!结果加完消费者发现 lag 还在涨,甚至反而触发了一堆 rebalance,消费更卡了。原因很简单:Kafka 的消费并行度上限不是消费者数量,而是分区数量。消费者数超过分区数之后,多出来的消费者就是纯闲置,等于白加。
这篇文章不讲空泛理论,直接围绕“消息堆积”这个故障场景,把排查路径、瓶颈判断、参数调优、Spring Boot 实战配置一次讲清楚。看完你应该能回答这几个问题:堆积时先看什么指标?为什么盲目加消费者没用?什么情况下加消费者有效?什么情况下应该改分区数、改消费逻辑、改提交方式?
1. 核心概念回顾:先搞清楚谁是瓶颈
在讨论“加消费者为什么没用”之前,先回顾 Kafka 消费模型里的三个关键概念:分区、消费者组、位移。
Kafka 一个 topic 可以拆成多个 partition,partition 是最小的并行单位。同一个消费者组内,一个 partition 最多只能被一个消费者实例消费。反过来说,一个消费者可以消费多个 partition。
这里就有一个关键公式:
消费者组最大并行度 = min(消费者实例数, 分区总数)消费者数超过分区数之后,超出的消费者分配的 partition 数量为 0,完全空闲。所以如果你只有一个分区,那不管你起 10 个消费者还是 100 个消费者,实际干活的消费者永远只有一个。
很多同学在 Spring Boot 里把concurrency从 3 调到 10,发现消费速度没变化,基本就是分区数小于 concurrency 导致的。
再来看消费吞吐的另外一个公式:
消费者组整体消费速度 = 单分区消费速度 × 分区数注意,这里是单分区消费速度,不是单消费者消费速度。一个消费者如果分配了多个分区,它处理这些分区是轮询拉取的,整体吞吐取决于每个分区消费的快慢。所以决定整体消费速度的两个变量是:
- 分区数
- 单个分区的消费耗时
如果你要提升消费速度,要么增加分区数(提高并行度上限),要么降低每一条消息的消费耗时(优化消费逻辑)。盲目加消费者,本质上没有动这两个变量中的任何一个,当然没用。
2. 为什么“盲目加消费者”没用
单纯加消费者无效,通常有四种典型情况。
2.1 分区数已经是瓶颈
这是最典型的情况。假设 topic 有 6 个分区,当前消费者组里有 10 个消费者,其中 4 个消费者没有分配到任何分区。这时候再加消费者,依然只有 6 个消费者在消费,新增消费者全部闲置。
判断方法也很简单:用工具查看消费者组的成员列表和分区分配情况。如果出现consumer-7的Current-offset和Log-end-offset全为 0,或者分配的分区列表为空,那基本就是分区数不够了。
这种情况的正确做法是扩容分区,而不是加消费者。扩容命令可以参考:
# 将 topic 分区数扩展到 12,需要根据业务实际情况评估 kafka-topics.sh --bootstrap-server localhost:9092 \ --alter \ --topic your-topic \ --partitions 12但是扩容分区有副作用,后面会单独讲。
2.2 单条消费逻辑太慢
如果每条消息消费耗时是 200ms,分区数假设是 10,那理论上单分区每秒只能处理 5 条消息,整个组每秒只能处理 50 条。这种情况下加消费者虽然能把吞吐往分区数上限靠,但只要消费逻辑不改,加多少消费者都逃不出“单分区消费速度 × 分区数”的天花板。
尤其是当消费逻辑涉及远程调用、数据库写入、外部 API 请求时,瓶颈往往就在这些同步等待上。消费者线程大部分时间都阻塞在 IO 等待上,CPU 几乎不干活。
2.3 下游系统扛不住
Kafka 消费堆积,不一定就是 Kafka 的问题,很可能是下游扛不住。比如消费后写 MySQL,如果数据库连接池打满、SQL 执行慢、锁等待严重,那消费者消费一条消息就要等很久。
这时候加消费者只会让更多请求打到下游,下游响应更慢,消费者端超时更多,堆积更严重。
2.4 频繁 Rebalance 导致消费停滞
加消费者不是改一个数字那么简单,它可能触发消费者组的 rebalance。Rebalance 期间,整个消费者组的所有消费者都会停止消费,等待重新分配分区。如果业务代码里max.poll.interval.ms设置太短,消费者处理一批消息耗时长,还没来得及发起下一次 poll,就被认为“已经死亡”,触发 rebalance。
如果频繁 rebalance,那消费者组会有大量时间处于“停止消费”的状态。看起来消费者数量很多,但其实都在等分配、等恢复,实际消费效率非常低。这个时候再加消费者,只会让 rebalance 更频繁。
3. 消息堆积的排查路径:从指标到日志
消息堆积出现时,第一件事不是改代码,而是确认“哪里最慢”。下面这套思路可以帮你快速定位。
3.1 先看消费 Lag
Lag 是衡量 Kafka 堆积最直接的指标。它表示消费者当前消费到的位移与最新消息位移之间的差值。
查看消费组 lag 的常用命令:
# 查看消费者组当前的消费位移和 lag kafka-consumer-groups.sh --bootstrap-server localhost:9092 \ --group your-consumer-group \ --describe输出大概长这样:
GROUP TOPIC PARTITION CURRENT-OFFSET LOG-END-OFFSET LAG CONSUMER-ID your-group your-topic 0 1000 5000 4000 consumer-1 your-group your-topic 1 800 5000 4200 consumer-2 your-group your-topic 2 0 5000 5000 consumer-3重点看两点:
- 每个 partition 的 lag 是否均匀
- 哪个消费者分配了哪个分区
如果某个 partition 的 lag 特别高,而其他 partition 已经追平,说明可能是 key 分布不均,或者某个分区被一个慢消费者拖住了。
3.2 看消费者线程是否活跃
用 jstack 抓一下消费者进程的线程栈,看消费者线程到底卡在哪个调用上:
# 找到 Java 进程 PID jps -l # 抓取线程栈 jstack <pid> > thread_dump.log消费者线程一般命名是consumer-<group>-<id>。看这些线程是处于RUNNABLE、WAITING还是BLOCKED状态。如果大量线程卡在数据库连接等待或者网络 IO 等待上,说明消费逻辑是瓶颈。
3.3 看 Kafka 监控指标
如果公司有 Kafka 监控面板,推荐重点盯这几个指标:
kafka.consumer:type=consumer-fetch-manager-metrics,client-id=*的records-lag-maxrecords-consumed-rate:每秒消费记录数bytes-consumed-rate:每秒消费字节数fetch-rate:每秒 fetch 请求次数fetch-latency-avg:fetch 平均延迟
如果fetch-rate很低但records-lag-max很高,说明消费者拉取频率太低,或者拉取间隔太长。
3.4 确认消息体大小
消息体大小对消费吞吐有很大影响。一条 1KB 的消息和一条 1MB 的消息,消费耗时完全不一样。如果生产端写入了大消息,消费者的网络 IO、反序列化都会变慢。
可以在消费端打印消息大小日志:
consumerRecord.value().toString().getBytes(StandardCharsets.UTF_8).length或者直接看 broker 端指标kafka.server:type=BrokerTopicMetrics,name=BytesInPerSec。
4. 定位瓶颈的四个检查点
消息堆积的瓶颈通常分布在四个位置:消费者、Broker、下游系统、网络。
4.1 消费者端
需要确认:
- 消费者数量是否已经等于分区数
- 消费逻辑单条耗时是多少
- 是否有大量 rebalance
- 消费线程是否频繁 GC
如果消费者线程频繁 Full GC,也会出现消费停顿。Kafka 客户端在 GC 停顿期间无法发送心跳,超过session.timeout.ms就会被踢出消费者组,触发 rebalance,导致消费进一步停滞。
4.2 Broker 端
消费者消费慢,不一定是消费者的问题,也可能是 Broker 响应慢。Broker 的磁盘 IO、网络带宽、分区副本同步情况都会影响 fetch 请求的响应速度。
建议关注 Broker 端指标:
kafka.network:type=RequestMetrics,name=TotalTimeMs,request=FetchConsumerkafka.server:type=KafkaRequestHandlerPool,name=RequestHandlerAvgIdlePercent
如果RequestHandlerAvgIdlePercent长期接近 0,说明 Broker 的请求处理线程已经满负荷,需要评估 Broker 扩容。
4.3 下游系统
很多消息堆积,根因在下游。比如消费后要写 Elasticsearch、MySQL、Redis、调用外部接口,下游慢,消费就会被拖住。
这里最容易被忽视的是数据库连接池大小和响应时间。如果连接池只有 10 个连接,而消费者线程有 20 个,那大量线程会阻塞在等待连接上。
排查时可以在消费逻辑里加耗时日志:
long start = System.currentTimeMillis(); // 消费逻辑 long cost = System.currentTimeMillis() - start; if (cost > 100) { log.warn("consume cost too much, partition={}, offset={}, cost={}ms", record.partition(), record.offset(), cost); }通过日志快速找到慢在哪一行调用。
4.4 网络与序列化
消费速度还受网络带宽和序列化方式影响。如果消息比较大,或者消费端反序列化比较慢,消费速率也会下降。比如使用 JSON 反序列化大量嵌套对象,在高吞吐场景下性能会明显不如 Protobuf 或 Avro。
5. 针对性调优手段
定位到瓶颈之后,再来决定怎么调。下面按瓶颈类型给出对应的调优方向。
5.1 分区数确实不够:扩容分区
如果分区数小于消费者数,且单分区消费速度已经很高,那就需要增加分区数。
但是要注意,扩容分区有代价:
- 现有消息不会自动重新分布,新分区是从当前时间点开始接收新消息
- 如果生产端使用 key 进行分区,扩容后 key 到 partition 的映射关系会变化,可能影响消息顺序
- 分区数只能增加,不能减少
所以扩容前要评估 topic 的分区数是否合理。一般建议分区数大于等于消费者最大并发数,保证消费者数有扩展空间。
kafka-topics.sh --bootstrap-server localhost:9092 \ --alter \ --topic order-events \ --partitions 24修改后重启消费者应用,让消费者组重新分配分区。
5.2 单分区消费速度慢:优化消费逻辑
这是最值得投入的方向。常见的优化手段包括:
批量处理替代单条处理
不要在消费者里一条条处理,把一批消息攒起来,批量写入下游。Kafka 消费者本身就支持一次拉取多条消息,配合 Spring Boot 的List<ConsumerRecord>批量监听,可以明显降低 IO 次数。
异步化非核心处理
如果消费逻辑里有非核心操作,比如发送通知、记录日志、做二次加工,可以丢到线程池异步执行。但要小心:异步处理后,消息已经提交了 offset,如果异步逻辑失败,消息就会丢失,需要自己保证可靠投递。
使用批量写库
对 MySQL、Elasticsearch 等存储的写入,尽量使用批量接口。单条 insert 和批量 insert 的性能差距可能是一个数量级。
5.3 Rebalance 频繁:调整消费者参数
很多堆积问题源于消费者被频繁踢出组。关键参数是max.poll.interval.ms。
这个参数表示消费者最多间隔多久发起一次 poll 请求。如果处理一批消息的时间超过这个值,消费者就会被判定为“死掉”,触发 rebalance。
假设你设置了max.poll.records=500,每条消息处理耗时 10ms,那一批就是 5 秒。如果max.poll.interval.ms还是默认的 300000ms,问题不大。但如果每批处理耗时超过 5 分钟,就要注意了。
建议根据实际消费耗时调整参数:
max.poll.records:控制单次拉取的消息数量max.poll.interval.ms:控制两次 poll 之间的最大间隔session.timeout.ms:控制会话超时时间heartbeat.interval.ms:控制心跳发送间隔,一般设为session.timeout.ms的三分之一
6. Spring Boot 集成 Kafka 的实战配置示例
下面给出一套生产环境更稳妥的 Spring Boot Kafka 消费配置。
6.1 application.yml 配置
spring: kafka: bootstrap-servers: kafka-1:9092,kafka-2:9092,kafka-3:9092 consumer: group-id: order-consume-group auto-offset-reset: latest enable-auto-commit: false key-deserializer: org.apache.kafka.common.serialization.StringDeserializer value-deserializer: org.apache.kafka.common.serialization.StringDeserializer max-poll-records: 100 fetch-min-bytes: 1024 fetch-max-wait-ms: 1000 max-poll-interval-ms: 300000 session-timeout-ms: 10000 heartbeat-interval-ms: 3000 listener: type: batch ack-mode: manual_immediate concurrency: 3几个参数说明:
enable-auto-commit: false:关闭自动提交,改为手动提交,避免消费逻辑失败导致 offset 丢失listener.type: batch:开启批量消费,一次拉取多条消息ack-mode: manual_immediate:手动提交,消费成功后立即 ackconcurrency: 3:消费者线程数,需要根据分区数调整,最大不要超过分区数
6.2 批量消费监听代码
import lombok.extern.slf4j.Slf4j; import org.apache.kafka.clients.consumer.ConsumerRecord; import org.springframework.kafka.annotation.KafkaListener; import org.springframework.kafka.support.Acknowledgment; import org.springframework.stereotype.Component; import java.util.List; @Slf4j @Component public class OrderMessageConsumer { @KafkaListener(topics = "order-events", groupId = "order-consume-group") public void onBatchMessage(List<ConsumerRecord<String, String>> records, Acknowledgment ack) { long start = System.currentTimeMillis(); try { for (ConsumerRecord<String, String> record : records) { // 这里是消费逻辑,尽量保证单条消费速度快 // 如果是写数据库,优先批量处理 processOrderEvent(record.value()); } // 手动提交 offset ack.acknowledge(); } catch (Exception e) { // 记录失败详情,进入补偿或重试机制 log.error("consume order event error, batch size={}", records.size(), e); // 这里根据业务决定是提交还是让消息重新消费 // 如果直接 ack,消息会丢失 // 如果一直不 ack,会阻塞消费进度,需要配合死信队列 } long cost = System.currentTimeMillis() - start; if (cost > 500) { log.warn("consume batch too slow, size={}, cost={}ms", records.size(), cost); } } private void processOrderEvent(String message) { // 模拟业务处理 // 真正的项目里这里可能是 ES 写入、MySQL 更新、外部接口调用 } }注意:批量消费时,如果中间一条消息处理失败,需要根据业务选择是整体跳过还是记录失败后继续。否则极端情况下会因为频繁重试导致消费者卡死。
6.3 手动提交的取舍
手动提交 offset 有三种常用方式:
- 消费完一批再提交:吞吐最高,但失败会丢消息
- 消费完一批先处理,再提交,失败就重试:可靠性高,但需要控制重试次数
- 每条消息处理成功后分别提交:可靠性最高,但性能最差
实际项目中推荐“批量处理 + 失败记录到死信 topic + 正常提交”的方案,兼顾吞吐和可靠性。
7. 消息堆积排查完整清单
建议把下面这份清单打印出来,遇到堆积问题时按顺序排查。
7.1 检查消费者状态
# 查看消费者组成员和 lag kafka-consumer-groups.sh --bootstrap-server localhost:9092 \ --group your-group \ --describe # 查看 topic 分区分布 kafka-topics.sh --bootstrap-server localhost:9092 \ --describe \ --topic your-topic确认:消费者数、分区数、每个消费者的分区分配、每个分区的 lag。
7.2 检查消费耗时
在消费逻辑里加耗时日志,或者用 Arthas 等工具监控方法耗时。确认慢在下游调用还是本地逻辑。
7.3 检查 Rebalance 频率
查看日志或者监控指标中的 rebalance 次数。如果一段时间内 rebalance 次数很多,重点检查max.poll.interval.ms是否合理,以及消费者是否频繁 GC 或 OOM。
7.4 检查下游
确认下游系统的连接池、写入 QPS、响应时间。在消费逻辑中继续调用下游前,先自己压测下游的极限吞吐。
7.5 检查消息大小和序列化
用工具确认消息的平均大小。如果单条消息过大,考虑在生产端做压缩,或者调整消费端参数。
8. 常见问题与排查方法
| 问题现象 | 可能原因 | 排查方式 | 解决方案 |
|---|---|---|---|
| 加消费者后 lag 不降 | 分区数已等于消费者数 | 查看消费者组 describe,确认空闲消费者 | 扩容分区,而不是继续加消费者 |
| lag 一直涨但消费者 CPU 很低 | 消费逻辑阻塞在 IO 或下游调用 | jstack 抓线程栈,看 BLOCKED/WAITING 状态 | 优化下游调用,改为异步或批量 |
| 频繁 rebalance | max.poll.interval.ms 过短或消费耗时过长 | 查看 rebalance 日志、检查 poll 间隔 | 调整 max.poll.interval.ms / max.poll.records |
| 消费重启后重复消费 | enable.auto.commit=true 且没有手动提交 | 检查 offset 配置 | 改用手动提交,保证处理成功后再 ack |
| 单一分区 lag 特别高 | key 分布不均导致热点分区 | 查看各分区 lag 分布 | 优化分区 key 设计,或者增加随机前缀 |
| 消费端报 DeserializationException | 消息格式变化 | 查看异常堆栈和日志 | 检查序列化配置,必要时加版本兼容 |
| fetch 请求超时 | Broker 负载过高或网络问题 | 查看 Broker 的请求处理耗时 | 优化 Broker 配置或扩容 Broker |
| 消费者线程数很多但吞吐低 | 大部分线程没有分配到分区 | 查看消费者组分配详情 | 减少 concurrency,或扩容分区 |
| 消费逻辑写库慢 | 数据库连接池打满或 SQL 慢 | 查看数据库慢查询日志 | 改用批量写入,或提升连接池大小 |
9. 最佳实践与调优建议
下面这些建议来自常见生产实践,不一定每条都符合你的场景,但方向基本通用。
9.1 分区数规划要留余量
新创建 topic 时,分区数要结合未来 1 到 2 年的数据量来评估。建议分区数是消费者最大并发数的 1.5 到 2 倍,留出扩消费者的空间。
9.2 第一次调优先小步验证
不要一次改一堆参数。先加日志和监控,观察当前消费速度,然后一次只改一个变量,对比效果。
9.3 消费逻辑要“快进快出”
消费者线程尽量只做消息转换和投递,把耗时操作放到异步线程池或下游任务队列中。如果必须同步调用下游,要设置明确的超时时间和重试机制,防止下游故障拖垮消费线程。
9.4 手动提交 offset 并配合死信队列
生产环境一定要关闭enable.auto.commit,使用手动提交。消费失败的消息写入专门的重试或死信 topic,不要一直阻塞主消费链路。
9.5 消息堆积要区分“积压”和“延迟”
有些场景下消费者速度正常,但生产者短时间写入量巨大,导致 lag 短暂升高。这种“积压”可以通过削峰填谷解决,不一定需要扩容消费者。真正需要关注的是持续增长的 lag,说明消费速度长期赶不上生产速度。
9.6 监控告警要分层
- lag 超过阈值告警
- 单分区 lag 超过阈值告警
- rebalance 次数异常增多告警
- 消费组长时间无消费告警
只有这些指标完整,才能在堆积刚出现时快速响应,而不是等到事故扩大后才排查。
10. 总结:加消费者的正确姿势
回到标题问题:Kafka 消息堆积,盲目加消费者确实没用。
正确的姿势是:先通过kafka-consumer-groups.sh --describe查看当前消费者组的 lag 分布和分区分配,再结合消费耗时、下游响应、rebalance 频率判断瓶颈位置。如果分区数已经是并发上限,就扩容分区;如果消费逻辑太慢,就优化消费逻辑;如果 rebalance 频繁,就调整消费者参数;如果下游扛不住,就先给下游减负。
加消费者这件事本身没错,但要在确认消费者数小于分区数、且消费逻辑不是瓶颈的前提下加。加完还要观察 lag 变化,确认有效果再继续加。
这套排查思路同样适用于 Kafka 面试题里常见的“消息堆积如何解决”场景。面试官真正想听的通常不是“加消费者”这个答案,而是你能说出分区数、消费者数、消费速度之间的关系,以及如何通过监控和日志定位瓶颈。
建议收藏备用,下次遇到 Kafka 消费堆积时照着排查一遍。