各位做后端开发的同学,应该都遇到过 Kafka 消息堆积的场景。消费 Lag 一路飙升,磁盘告警,业务数据迟迟不更新,这时候很多人第一反应就是“加消费者,水平扩容嘛”。但实际加完却发现,消费者数量加了不少,堆积反而纹丝不动,甚至有的消费者节点根本就没在工作。
这篇文章就围绕“Kafka 消息堆积,为什么不能盲目加消费者”这个核心问题展开,从 Kafka 的消费模型讲起,梳理消息堆积的真实原因,再给出一套完整的排查思路和解决方案。无论你是刚入门 Kafka 的新手,还是已经被线上堆积问题折磨过的开发,这篇文章都值得收藏备用。
1. 消息堆积的本质:先理解 Kafka 消费模型
1.1 什么是 Kafka 消息堆积
Kafka 是基于发布订阅模式的消息中间件,生产者将消息写入 Topic,消费者从 Topic 拉取消息进行业务处理。所谓“消息堆积”,简单来说就是生产者的写入速度长期大于消费者的处理速度,导致消息在 Kafka Broker 上越积越多。
在 Kafka 中,每个 Topic 会被划分为多个 Partition(分区),消息按分区存储。消费者通过记录 Offset(偏移量)来标记自己消费到哪一条消息。当消费者的 Offset 和 Partition 最新的 Offset 差距越来越大时,就说明消息在堆积。
这里有一个重要概念:消费 Lag(消费落后值)。Lag = 分区最新 Offset - 当前消费 Offset。Lag 越大,堆积越严重。Kafka 本身不会主动删除未消费的消息,只会根据日志保留策略(retention)定期清理过期数据,所以堆积的消息不会自动消失,要么被消费掉,要么等到过期被删除。
1.2 分区与消费者的关系
理解 Kafka 消息堆积,绕不开分区(Partition)和消费者(Consumer)的关系。Kafka 的消息模型有几个关键规则:
第一,一个分区在同一时刻只能被同一个消费组(Consumer Group)内的一个消费者线程消费。这是 Kafka 保证分区内消息有序性的基础。
第二,一个消费者可以同时消费多个分区。消费者和分区之间是多对多的关系,但约束在每个分区只会被组内一个消费者持有。
第三,当消费者数量大于分区数量时,必然有消费者分配不到任何分区,处于空闲状态。
举个例子,假设有一个 Topic,它有 3 个分区,消费组内有 5 个消费者实例。那么这 5 个消费者中只会分配 3 个去消费分区,另外 2 个消费者完全闲置,不处理任何消息。
这就是盲目加消费者没用的根本原因所在:如果 Topic 的分区数不增加,消费者加再多,也只是增加闲置的消费者实例,消费能力并不会提升。很多人忽略了这个前置条件,一看到堆积就扩容消费者节点,结果只是白白浪费机器资源。
1.3 加消费者的正确前提
那什么时候加消费者是有效的?答案是:当前消费者数量小于分区数。
假设 Topic 有 10 个分区,当前消费组只有 2 个消费者,每个消费者平均要消费 5 个分区。此时把消费者扩展到 5 个,每个消费者平均只消费 2 个分区,单分区消费压力变小,整体消费速度自然提升。
但如果 Topic 只有 3 个分区,当前已经有 3 个消费者,再加到 10 个消费者,实际依然只有 3 个消费者在工作,其余 7 个都在空转。
所以在扩容消费者之前,第一件事应该是确认 Topic 的分区数。
2. 消息堆积的六大常见根因
既然不能一上来就加消费者,那 Kafka 消息堆积通常由哪些原因引起?我在实际项目里总结下来,主要有六类。
2.1 上游生产速度超过下游消费速度
这是最直白的堆积原因。比如大促秒杀场景,突然涌入大量订单消息,生产者瞬间写入了百万条消息,而消费者的处理逻辑需要查数据库、调外部接口,单条消息处理耗时在几百毫秒甚至秒级,消费速度远远跟不上生产速度,Lag 就会快速上升。
这种堆积是短时流量冲击造成的,也可能是因为系统长期处于“生产者写入快、消费者处理慢”的状态。前者属于正常流量波动,后者属于设计缺陷。
2.2 分区数不足导致并行度受限
这是“加消费者没用”最常见的场景。Topic 创建时分区数设置过小,例如只设置了 3 个分区,但业务量已经增长到需要 30 个消费者并行处理。这时消费者侧无论怎么扩容,实际并行度只有 3,消息积压就无法消化。
很多团队创建 Topic 时为了省事,直接用了默认分区数(如 1 个或 3 个),后续业务增长时又没有及时评估分区数是否满足消费并行度需求,最终导致堆积。
2.3 消费者处理逻辑耗时过高
消费者拉取到消息后,需要执行业务逻辑。如果消费逻辑中存在慢 SQL、外部 API 调用超时、循环嵌套、大对象序列化等耗时操作,单条消息的处理时间会大大增加。
Kafka 消费者是拉取模型,本地会有一个 poll 循环。每次 poll 拉取一批消息,然后由业务线程处理。如果处理时间太长,下一轮 poll 就会延后。更严重的是,如果处理时间超过了max.poll.interval.ms(默认 5 分钟),消费者会被认为“失联”,触发 Rebalance,分区被分配给其他消费者,造成重复消费和更大的消费延迟。
2.4 Offset 提交方式不合理
消费者消费完消息后,需要提交 Offset,告诉 Kafka“这条消息我已经处理完了”。 Offset 提交分为自动提交和手动提交,两种方式都有各自的坑。
自动提交模式下,消费者定期提交当前拉取到的 Offset,而不是处理完的 Offset。如果消费者拉取了一批消息,还没处理完就到了自动提交时间点,进程突然宕机,重启后会从已提交的 Offset 继续拉取,这部分消息就丢失了。反过来说,如果业务逻辑处理完后程序崩溃,但 Offset 已经在此之前提交了,重启后就会跳过一批消息。
手动提交模式下,如果业务代码处理完消息后忘记提交 Offset,或者提交逻辑写在了异常路径之外,那么每次重启后都会从旧的 Offset 开始重新消费。更常见的是,提交时机不对——例如在调用consumer.poll()之后就立刻提交 Offset,而不是在处理完这批消息之后再提交,极端情况下会丢消息。虽然丢消息不算严格意义上的堆积,但会导致业务数据不一致,从用户视角看,消息“堆积”在那里永远处理不完。
2.5 频繁 Rebalance 导致消费停滞
Kafka 消费组内出现成员变化(消费者加入、离开、崩溃)时,会触发 Rebalance(再平衡),将分区在消费者之间重新分配。Rebalance 期间,所有消费者都会停止消费,分区无法被处理,相当于整个消费组暂停服务。
以下几种情况容易频繁触发 Rebalance:
- 消费者处理消息耗时超过
max.poll.interval.ms,被 Kafka 判定为失败并踢出消费组。 - 消费者与 Broker 之间的心跳超时,
session.timeout.ms设置过短,网络抖动导致误判。 - 消费者进程频繁 Full GC,导致线程长时间暂停,无法发送心跳。
- 消费者实例频繁重启或扩容缩容,每次变化都会触发全量 Rebalance。
- 消费线程在处理消息时抛出未捕获异常,导致消费者退出。
每次 Rebalance 都会中断消费,造成 Lag 上升。如果 Rebalance 频繁发生,消费组大部分时间都花在分区重新分配上,堆积问题会越来越严重。
2.6 下游依赖成为瓶颈
消费者的速度往往不取决于自身,而取决于下游依赖。比如消费消息时需要写入数据库,如果数据库连接池打满、表锁竞争严重,即使消费者逻辑写得再高效,消息也会在等待数据库响应中积压。同样,如果消费者需要调用第三方 HTTP 接口,而第三方接口响应很慢,消费速度也会被拖慢。
这类问题的典型特征是:Kafka 侧消费 Lag 很高,但消费者所在机器的 CPU、内存利用率都不高,线程大量阻塞在外部 IO 等待上。
3. 环境准备与版本说明
在动手排查之前,先明确一下文章使用的环境。Kafka 版本差异会导致命令和参数有所不同,但核心思路是一致的。
本文示例以常见环境为例:
- Kafka 版本:2.8 及以上,采用 ZooKeeper 或 KRaft 模式均可。
- Java 版本:JDK 8 / JDK 11。
- 消息客户端:Spring Kafka 2.8 或 Kafka Client 3.x。
- 构建工具:Maven。
- 操作系统:Linux / macOS。
如果你的项目使用的是 Kafka 1.x 或 2.x 早期版本,部分命令行参数会略有差异,以实际环境帮助文档为准。重点演示的是排查思路和配置逻辑,不同版本下这些原理是通用的。
4. 消息堆积的完整排查流程
遇到消息堆积,不要急着改代码,先做排查。下面这套流程可以帮助你准确定位堆积根因。
4.1 查看消费组与消费 Lag
第一步是确认堆积到底有多严重,以及堆积发生在哪些分区。
使用 Kafka 自带的命令行工具即可查看消费组详情:
# 查看消费组列表 kafka-consumer-groups.sh --bootstrap-server localhost:9092 --list # 查看指定消费组的消费进度和 Lag kafka-consumer-groups.sh --bootstrap-server localhost:9092 --describe --group order-service-group执行结果类似:
GROUP TOPIC PARTITION CURRENT-OFFSET LOG-END-OFFSET LAG CONSUMER-ID order-service-group order-topic 0 1000 5000 4000 consumer-1 order-service-group order-topic 1 2000 9000 7000 consumer-2 order-service-group order-topic 2 1500 8000 6500 consumer-3其中:
CURRENT-OFFSET:当前消费组已经消费到的 Offset。LOG-END-OFFSET:分区中最新的消息 Offset。LAG:还未消费的消息条数。CONSUMER-ID:当前正在消费该分区的消费者实例。
注意观察各分区的 LAG 分布是否均匀。如果 LAG 集中在某一个或某几个分区,说明分区分配不均匀,或者某个消费者处理能力较弱;如果所有分区的 LAG 都很高,说明整体消费速度都跟不上。
4.2 确认消费者数量与分区数量
这是判断“加消费者是否有用”的关键一步。
统计 Topic 的分区数:
# 查看 Topic 详情,包括分区数、副本数等 kafka-topics.sh --bootstrap-server localhost:9092 --describe --topic order-topic输出中会列出所有分区信息,例如Topic: order-topic PartitionCount: 3。
再查看消费组内有多少个消费者实例。可以数一下上一步--describe输出中的CONSUMER-ID数量,或者在消费者日志中查看。如果消费者实例数 >= 分区数,说明加消费者大概率没用,瓶颈不在消费者数量上。
4.3 查看消费者分配情况
确认消费者数量后,进一步查看每个消费者分配了哪些分区。
继续使用kafka-consumer-groups.sh:
kafka-consumer-groups.sh --bootstrap-server localhost:9092 --describe --group order-service-group --members输出示例:
CONSUMER-ID HOST ASSIGNMENT consumer-1 /10.0.0.1 order-topic-0, order-topic-1 consumer-2 /10.0.0.2 order-topic-2 consumer-3 /10.0.0.3 (无分区)从这里可以清楚看到,consumer-3 没有分配到任何分区,它就是一个闲置消费者。此时哪怕继续加 consumer-4、consumer-5,也都是同样闲置,根本不会提升消费能力。
4.4 定位消费慢的具体环节
确认分区数和消费者数量都没问题后,就要判断消费者为什么慢。这一步通常需要在消费者机器上观察指标。
重点查看以下几项:
第一,消费者所在的 JVM 进程是否频繁发生 Full GC。Full GC 会导致线程长时间停顿,无法消费消息。可以使用jstat命令观察:
# 每 1 秒输出一次 GC 情况,共输出 10 次 jstat -gcutil <pid> 1000 10如果 FGC 列的数字持续增长,且 FGCT 时间不断增加,说明 JVM GC 已经是瓶颈。
第二,消费者线程是否大量阻塞。用jstack导出线程快照,查看消费者线程处于什么状态:
jstack <pid> > jstack.log重点看consumeMessage相关的线程是否大量处于WAITING或BLOCKED状态。如果线程阻塞在数据库连接上,说明数据库是瓶颈;如果阻塞在SocketRead,说明外部接口调用慢。
第三,观察机器的基础监控指标。CPU 使用率、内存使用率、磁盘 IO、网络带宽,每一项都可能是瓶颈。例如 CPU 打满说明消费逻辑中密集计算太多;磁盘 IO 高说明消息值过大导致序列化开销高;网络带宽打满说明消息内容太大。
4.5 检查 Rebalance 频率
如果消费者经常“掉线”又被重新分配分区,消费进度会反复回退,Lag 也会居高不下。
在消费者日志中搜索关键字,例如Rebalance、rebalance、Assignments、Group coordinator。正常情况下,消费组启动后不应该频繁出现 Rebalance 日志。如果短时间内出现多次分区重新分配,需要进一步排查心跳超时和max.poll.interval.ms配置。
同时观察消费者与 Broker 之间的网络状况,心跳超时往往是网络抖动造成的,也可能是session.timeout.ms设置过短,稍微一点网络波动就触发了超时判定。
5. 解决 Kafka 消息堆积的有效方案
定位到具体原因后,就可以对症下药了。下面这些方案按场景分类,你可以根据排查结果选择组合使用。
5.1 方案一:扩展 Topic 分区数
当确认“分区数不足导致消费并行度不够”时,扩容分区是最直接的手段。
# 将 order-topic 扩展到 12 个分区 kafka-topics.sh --bootstrap-server localhost:9092 --alter --topic order-topic --partitions 12这里必须强调一个重要问题:Kafka 分区数只能增加,不能减少。扩容前一定要评估清楚,分区数设置得过大,会增加 Broker 的元数据管理开销和文件句柄占用,也会增加同一消费组内的管理成本,但通常来说适度的分区扩容是安全的。
分区扩容之后,还要确认消费者能随之扩展。假设原来 3 个分区对应 3 个消费者,扩到 12 个分区后,理想情况下消费组内应该扩展到 12 个消费者实例来充分利用分区并行度。如果消费者实例不增加,只是分区数变多,那么每个消费者要消费的分区数变多,单消费者压力反而增大,消费速度不一定能提升。
实际项目中,如果 Topic 已经存在大量消费者,而且消费者数量跟不上分区数,建议同时完成分区扩容和消费者实例扩容,并逐步重启消费组,避免一次性触发大规模 Rebalance。
5.2 方案二:优化消费者单条处理耗时
如果分区数和消费者数量都合理,但单条消息处理太慢,就需要从消费逻辑本身入手。
常见的优化手段包括:
- 精简消费逻辑,把耗时的非核心操作放到异步线程中执行。
- 尽量避免在消费者线程中调用慢速外部接口,可以考虑批量聚合后统一调用。
- 使用批量处理,一次性处理多条消息,而不是一条条处理。
Kafka 消费者天然支持批量拉取,关键在于业务代码怎么写。下面是一个优化示例,处理消息时不单条入库,而是攒一批后批量写入:
// 文件路径:src/main/java/com/example/kafka/BatchConsumer.java @Component public class BatchConsumer { private static final Logger log = LoggerFactory.getLogger(BatchConsumer.class); private static final int BATCH_SIZE = 100; @KafkaListener(topics = "order-topic", groupId = "order-service-group") public void onMessage(List<ConsumerRecord<String, String>> records) { List<OrderMessage> orderList = new ArrayList<>(records.size()); for (ConsumerRecord<String, String> record : records) { try { OrderMessage message = JSON.parseObject(record.value(), OrderMessage.class); orderList.add(message); if (orderList.size() >= BATCH_SIZE) { batchInsert(orderList); orderList.clear(); } } catch (Exception e) { log.error("消息解析失败,record={}", record.value(), e); } } if (!orderList.isEmpty()) { batchInsert(orderList); } } private void batchInsert(List<OrderMessage> orderList) { // mybatis 批量插入或者其他批量写入逻辑 // orderMapper.batchInsert(orderList); } }使用批量处理时要注意,@KafkaListener接收 List 参数需要配置批量工厂和消费者配置,才能把多条消息一次性拉入监听方法中。
5.3 方案三:合理调整消费者核心参数
Kafka Client 提供了一些参数,直接影响消费速率和稳定性。以下参数需要重点关注。
max.poll.records:单次 poll 拉取的最大消息条数,默认 500。如果单条消息处理较慢,可以适当调小,比如 100 或 200,避免单次拉取过多消息导致处理时间超过max.poll.interval.ms。
示例配置:
spring.kafka.consumer.max-poll-records=200 spring.kafka.consumer.max-poll-interval-ms=300000 spring.kafka.consumer.session-timeout-ms=15000 spring.kafka.consumer.heartbeat-interval-ms=5000 spring.kafka.consumer.enable-auto-commit=false如果使用原生 Kafka Client,在Properties中配置:
// 文件路径:src/main/java/com/example/kafka/KafkaConsumerConfig.java Properties props = new Properties(); props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092"); props.put(ConsumerConfig.GROUP_ID_CONFIG, "order-service-group"); props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName()); props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName()); props.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, "false"); props.put(ConsumerConfig.MAX_POLL_RECORDS_CONFIG, "200"); props.put(ConsumerConfig.MAX_POLL_INTERVAL_MS_CONFIG, "300000"); props.put(ConsumerConfig.SESSION_TIMEOUT_MS_CONFIG, "15000"); props.put(ConsumerConfig.HEARTBEAT_INTERVAL_MS_CONFIG, "5000");需要特别提醒的是,session.timeout.ms和heartbeat.interval.ms要配合设置。一般推荐session.timeout.ms为heartbeat.interval.ms的 3 倍左右。如果设置过短,网络轻微波动就会触发 Rebalance。
手动提交 Offset 的代码示例:
// 文件路径:src/main/java/com/example/kafka/ManualCommitConsumer.java @Component public class ManualCommitConsumer { @KafkaListener(topics = "order-topic", groupId = "order-service-group") public void onMessage(ConsumerRecord<String, String> record, Acknowledgment ack) { try { // 处理业务逻辑 process(record.value()); // 处理成功后再提交 offset ack.acknowledge(); } catch (Exception e) { // 记录失败日志,根据业务决定是否重试 log.error("消息消费失败,topic={}, offset={}", record.topic(), record.offset(), e); // 不提交 offset,下次 poll 会继续拉取该消息 } } }手动提交 Offset 时,必须确保“先处理业务,后提交 offset”,否则会出现消息丢失。
5.4 方案四:异步化与多线程消费
当单消费者处理能力有限时,可以在消费者内部引入多线程,将拉取和处理解耦。
常见的模式是:消费者线程只负责拉取消息,将消息放入内存队列(如LinkedBlockingQueue)或线程池中,由工作线程并发处理。
示例代码如下:
// 文件路径:src/main/java/com/example/kafka/AsyncConsumer.java @Component public class AsyncConsumer { private static final int CORE_POOL_SIZE = 8; private static final int MAX_POOL_SIZE = 16; private static final int QUEUE_CAPACITY = 10000; private final ExecutorService executor = new ThreadPoolExecutor( CORE_POOL_SIZE, MAX_POOL_SIZE, 60L, TimeUnit.SECONDS, new ArrayBlockingQueue<>(QUEUE_CAPACITY), new ThreadFactoryBuilder().setNameFormat("kafka-consumer-worker-%d").build(), new ThreadPoolExecutor.CallerRunsPolicy() ); @KafkaListener(topics = "order-topic", groupId = "order-service-group") public void onMessage(ConsumerRecord<String, String> record) { executor.submit(() -> { try { process(record.value()); } catch (Exception e) { log.error("异步消费失败,offset={}", record.offset(), e); } }); } }注意,使用多线程消费会带来两个新问题。
第一,消息顺序可能无法保证。不同线程处理同一个分区的不同消息时,顺序可能颠倒。如果业务对消息顺序有严格要求,多线程方案不适合。
第二,失败重试和 Offset 提交变得复杂。由于消息已经提交到线程池,消费者线程无法感知处理结果。如果工作线程处理失败,不能简单地通过不提交 Offset 来重试,需要额外的失败重试队列或补偿机制。
因此,异步多线程方案适合对顺序不敏感、单条消息失败可以单独补偿的业务场景。
5.5 方案五:临时堆积时增加 Topic 分区
有时候流量突增只是暂时的,系统整体设计没有问题,只是短期 Lag 上涨。此时如果为了峰值流量长期保持大量消费者和分区数,成本不划算。
可以考虑在临时堆积期间手动扩容分区和消费者,流量恢复正常后再缩容消费者实例。注意,Kafka 分区数不能减少,所以“临时扩容分区”意味着分区数会永久增加,后续要接受这个副作用。实际操作时更稳妥的做法是:评估峰值流量的持续时间,如果持续时间长,扩容分区;如果只是短时峰值,可以临时扩展消费者实例数(前提是分区数允许)。
5.6 方案六:从源头降低生产速率影响
如果上游生产速率过高是常态,除了提高下游消费能力,还可以从 Topic 设计层面缓解。
一种做法是按照业务优先级拆分 Topic,把实时性要求高的消息和实时性要求低的消息分开。例如订单状态变更实时性要求高,单独用一个 Topic;用户行为日志允许一定延迟,放到另一个 Topic。避免大量低优先级消息挤占高优先级消息的消费资源。
另一种做法是对突发的流量做削峰填谷,生产者侧增加限流,或者将部分消息先写入临时存储,再异步落到 Kafka。这样可以减少 Broker 压力,但会增加架构复杂度,需要结合实际场景权衡。
6. 常见问题与排查速查表
把消息堆积排查过程中最常见的问题整理成一张速查表,方便你遇到类似情况时快速定位。
| 问题现象 | 常见原因 | 排查思路 | 解决方向 |
|---|---|---|---|
| 加消费者后 Lag 不减 | 消费者数量已大于等于分区数,新增消费者闲置 | 查看 Topic 分区数、消费组成员分配 | 扩容分区,或优化消费逻辑 |
| 单分区 LAG 特别高 | 分区分配不均匀,或某个消费者处理能力差 | 查看 --members 分配情况和消费者监控 | 调节分区分配策略,排查问题节点 |
| 消费者频繁掉线 | session.timeout.ms过短或网络不稳定 | 查看心跳日志和 Rebalance 日志 | 适当调大 session.timeout.ms,减少心跳间隔 |
| 处理时间超过 max.poll.interval.ms | 单条消息处理太慢,poll 循环阻塞 | 检查消费逻辑中的耗时操作 | 优化逻辑、调整为批量处理、多线程异步化 |
| Offset 不向前推进 | 手动提交 Offset 代码未执行或异常 | 检查提交逻辑是否在 try 块内 | 确保处理成功后再提交,异常时记录日志 |
| 消费端资源利用率很低但 Lag 很高 | 下游依赖成为瓶颈,线程阻塞在外部 IO | jstack 查看线程状态 | 优化数据库连接池、外部接口调用,加入缓存 |
| Topic 分区数过少 | 创建 Topic 时使用默认分区数 | 查看 Topic 分区数 | 适当增加分区数,注意分区只能增不能减 |
| 消息消费重复 | Rebalance 后重复拉取未提交 Offset 的消息 | 检查 Offset 提交时机和幂等处理 | 使用手动提交,业务侧做幂等 |
7. 最佳实践与工程建议
7.1 消费组和分区规划设计
创建 Topic 时就要考虑分区数。不要凭感觉拍脑袋,可以参考以下公式估算:
分区数 = 预期的目标消费速率(条/秒) / 单个分区的消费能力(条/秒)单个分区的消费能力很难精确估算,通常可以通过压测得到。但有一个原则是:分区数宁多勿少,因为 Kafka 的分区数只能增加不能减少。如果一开始设太少,后面扩容就是一次风险操作。
同时要规范 Topic 命名和消费组命名。推荐格式类似业务域.事件类型,例如order.created、user.login。消费组命名最好能和业务模块对应,这样排查问题时一眼就能看出来是哪条链路。
7.2 监控和告警体系
消息堆积问题最重要的是“早发现”。建议从以下几个方面建立监控:
第一,消费 Lag 监控。Kafka 提供了kafka-consumer-groups.sh命令行工具,可以配合脚本定时采集 Lag 数据。如果使用 Kafka 3.x,Lag 监控会集成到 Kafka 自身指标中。如果引入 Kafka Exporter 和 Prometheus,可以直接在 Grafana 中查看消费组 Lag 变化曲线。
第二,消费者 JVM 监控。重点关注 Full GC 频率、堆内存使用率、线程阻塞情况。Full GC 会影响消费者心跳,严重的会导致 Rebalance。
第三,下游依赖监控。数据库连接池使用率、外部接口响应时间、消息队列中的积压数量,这些指标的异常往往比 Kafka 本身的 Lag 更早暴露问题。
设置告警时要注意阈值合理性。Lag 有轻微波动是正常的,不要设得过于敏感。一般建议设置两级告警:一级是 Lag 超过某阈值且持续 5 分钟,提示关注;二级是 Lag 持续上涨且超过积压上限,触发紧急处理。
7.3 代码层面的稳定性建议
消费者代码最容易踩的坑是异常处理不完善。
第一,消费者线程中不要抛出未捕获的异常。未捕获异常会导致消费者消费线程退出,但进程仍然存活,测试环境里很难发现。建议在消费入口处统一捕获异常,记录日志,根据业务场景决定是重试还是丢弃。
第二,消费逻辑要支持幂等。Kafka 在异常和 Rebalance 场景下可能会重复投递消息,消费端如果不对重复消息做幂等处理,就会出现数据重复。常见做法是利用数据库唯一索引、Redis SETNX 或者业务号去重。
第三,消息处理失败要设计重试机制。最简单的是将失败消息写入一个重试 Topic,由独立的消费者做延迟重试;或者使用 Kafka 的RetryTopicConfiguration(Spring Kafka 提供)实现自动重试。
7.4 生产环境变更注意事项
涉及 Kafka Topic 分区扩容、消费组重置这类操作,以下几点必须重视。
扩容分区属于高危变更,Kafka 不支持缩容,操作前必须确认分区数增加的合理性,并且最好在业务低峰期执行。扩容后消费者可能需要触发一次 Rebalance 才能感知到新分区,要留意 Rebalance 对现有消费的影响。
修改消费者参数时,小步快跑,每次只改动一个参数,观察 Lag 和消费者稳定性。不要一次性修改多个参数,否则出了问题无法定位是哪个参数引起的。
如果需要对消息进行重放或者重置 Offset,不要在生产环境直接操作。先在测试环境验证逻辑,再在维护窗口操作,并且操作前备份消费组的 Offset 信息,便于快速回滚。
8. 总结
回到文章标题:Kafka 消息堆积,盲目加消费者为什么没用?核心原因就是 Kafka 消费模型的并行度上限由分区数决定,消费者数量超过分区数之后再多都是空转。所以在处理堆积问题时,正确的顺序是:先确认分区数和消费者的关系,再看消费 Lag 分布,然后排查消费者自身的处理瓶颈,最后才是扩容或优化代码。
Kafka 消息堆积本身并不可怕,可怕的是堆积发生后找不到根因,盲目操作反而让问题发酵。希望这篇文章能帮你建立一套完整的排查思路。如果你在实际项目中也遇到过类似的 Kafka 堆积问题,或者有更好的处理方案,欢迎在评论区交流。下篇文章可以继续聊聊 Kafka 的 Rebalance 原理和消费者分区分配策略,感兴趣的可以先收藏备用。