
线上 Kafka 消费积压告警又响了你的第一反应是什么很多人会直接说加机器、加消费者实例把分区分摊出去也就是“扩容”。如果这是你的第一反应那这篇文章建议认真读完。我见过不少团队在遇到积压时第一选择就是扩容结果机器加了三台积压不但没降下来反而因为频繁 rebalance、下游数据库压力暴涨把问题放得更大。Kafka 积压问题真正难的地方不是“怎么扩容”而是“怎么找到导致积压的瓶颈点”。扩容只是众多手段中的一种而且往往不是最优解。这篇文章会从 Kafka 积压的成因讲起分析“只会扩容”这种初学者思维的局限性再给出一套从指标排查、日志分析到代码优化的完整方案。文中的消费端示例基于 Java Spring Boot代码可以直接落地也可以改造成其他语言版本。1. Kafka 消息积压到底是怎么产生的1.1 什么是消息积压Kafka 是一个基于发布订阅模式的分布式消息队列。生产者把消息写入 topic 的某个分区partition消费者通过消费者组consumer group订阅 topic各自负责消费一部分分区。所谓“消息积压”可以简单理解成生产者写入消息的速度超过了消费者消费消息的速度导致 Kafka 中未被消费的消息越来越多。专业一点的描述是某个分区最新写入的 offset 与消费者当前已提交的 offset 之间的差值越来越大这个差值通常被称为“消费滞后量”英文叫 Consumer Lag。可以用下面这个思路来理解生产者写入位置指最新一条消息写入分区后的 offset。消费者消费位置指消费者处理完并提交的 offset。Lag 最新写入 offset − 已提交 offset。如果一条消息写入后消费者很快消费并提交Lag 会保持在一个很小的范围内。如果消费者处理太慢或者消费者线程阻塞、宕机Lag 就会不断上涨最终变成“积压”。1.2 消息积压会造成什么影响很多人觉得消息积压无非就是“消息晚点处理”影响不大。但在真实业务中积压往往意味着线上事故级别的风险。第一业务实时性受损。比如下单后需要发券、发短信、更新库存如果这些操作依赖 Kafka 异步消费积压会导致用户下单后迟迟收不到通知或者库存数据更新不及时进而出现超卖。第二消息存活时间有限。Kafka 有日志保留策略默认可能只保留几天数据。如果积压严重部分消息可能在消费者还没消费到之前就被清理掉造成数据丢失。第三故障扩散。消费积压时消费者线程长时间忙碌可能导致后续消息处理超时、数据库连接池被占满、下游系统被拖垮。这时候如果盲目扩容消费者反而会放大对下游的压力。2. 为什么说“只会扩容”是初学者的解决方式2.1 扩容解决的只是表象这里说的扩容指的是增加消费者实例数量或者提高单个消费者的并发线程数让更多分区可以被并行消费。表面上看消费者变多了消费总吞吐应该变大Lag 应该下降。但这个逻辑成立的前提是瓶颈真的在消费者自身。如果瓶颈在下游数据库比如每次消费都要执行一次 insert数据库单表写入能力有限那么你加再多的消费者数据库连接池和写入锁会成为新的瓶颈最终结果可能是数据库负载飙升。消费者线程大量阻塞等待数据库响应。消费速度没有明显提升。数据库出现慢查询甚至宕机。这种情况下盲目扩容不仅解决不了积压还会放大对下游系统的压力。真正该做的是减少不必要的数据库写入或者把单条 insert 改造成批量 insert。所以“初学者才会用扩容解决 Kafka 积压问题”这句话并不是说扩容完全没用而是说“无脑扩容”是初学者思维。成熟的思路是先定位瓶颈再选择对应的优化手段。扩容只是候选方案之一不是默认方案。2.2 扩容前必须回答的四个问题在决定扩容之前建议先回答下面四个问题。如果答不上来说明还没有找到根因扩容大概率是无效操作。第一个问题当前消费者数量与分区数的关系是什么Kafka 的基本模型是一个分区在同⼀个消费者组内同时只能被一个消费者实例消费。如果消费者实例数已经大于等于分区数那你继续增加消费者实例新增的实例也分不到任何分区扩容等于白扩。第二个问题单个消费者的消费速率是多少瓶颈在哪个环节是反序列化慢还是业务逻辑慢还是下游接口调用慢如果单条消息处理耗时 200ms其中 180ms 花在下游 HTTP 调用上那优化方向应该是减少调用次数或改成异步调用而不是加机器。第三个问题扩容后会不会引发 rebalance 风暴当消费者实例发生变更时Kafka 会触发 rebalance也就是消费者组重新分配分区。rebalance 期间整个消费组会停止消费。如果频繁扩容、缩容rebalance 会反复打断消费反而加剧积压。第四个问题扩容的容量公式是否算过单个消费者吞吐可以简单估算为单消费者吞吐条/秒 1000 / 单条处理耗时毫秒 × 并发线程数 消费组总吞吐 ≈ 单消费者吞吐 × min(消费者实例数, 分区数)举个例子如果单条消息处理耗时 200ms单消费者单线程每秒只能处理 5 条。即使你有 10 个消费者实例但 topic 只有 4 个分区实际总吞吐也就是 20 条/秒。这种情况下扩容消费者实例毫无意义正确的做法是增加分区数或者提升单条处理效率。3. 定位积压根因的三层排查法3.1 第一层看消费指标处理积压问题第一步不是改代码而是收集指标。核心指标有三个消费组 Lag、消费速率、分区 Lag 分布。Kafka 自带的命令行工具可以直接查看消费组详情。Kafka 2.x 和 3.x 的常见命令如下kafka-consumer-groups.sh \ --bootstrap-server localhost:9092 \ --group order-consumer-group \ --describe输出中会包含每个分区的CURRENT-OFFSET当前已提交消费位置。LOG-END-OFFSET分区最新写入位置。LAG两者差值。如果发现所有分区的 LAG 都在快速增长说明消费者整体消费速度跟不上。如果只有某个分区 LAG 很大其他分区正常那可能是分区数据倾斜或者该分区对应的消费者实例出现异常。生产环境建议把 Lag 接入监控系统例如 Prometheus Grafana或者使用 Kafka 自带的 JMX 指标。设置一个合理的告警阈值比如 Lag 超过 10000 持续 5 分钟就告警比每天手动看命令输出要可靠得多。3.2 第二层看消费者日志指标能告诉你“积压了”但不会告诉你“为什么积压”。下一步需要看消费者日志。重点观察以下几类日志rebalance 相关日志如果日志中频繁出现Rebalance started、Rebalance completed说明消费者组在反复发生分区重分配。常见原因是消费者处理超时触发了max.poll.interval.ms限制导致消费者被认为已经失联。消费耗时日志如果业务代码里记录了每条消息的处理耗时可以通过日志聚合分析单条消息平均耗时、P99 耗时。异常与重试日志比如数据库连接超时、外部接口返回 5xx、反序列化失败等这些异常会导致消费线程反复重试直接拉低消费速率。在很多实际案例中积压的根因并不是消费速度慢而是消费者频繁 rebalance。只要 rebalance 一发生整个消费者组会暂停消费Lag 自然快速上涨。这种情况加机器只会让 rebalance 更频繁积压更严重。3.3 第三层看业务链路耗时如果消费日志没有明显异常单条消息处理耗时却很高就需要借助链路追踪工具来分析耗时分布。常见的方案是 SkyWalking、Zipkin 或 Micrometer Tracing。核心思路是把一次消费拆成多个阶段分别统计耗时。假设一次消费包含以下几个步骤从 Kafka 拉取消息。反序列化消息体。查询数据库获取关联数据。调用外部积分服务。组装结果并写入数据库。通过链路追踪能直接看到 200ms 到底花在哪一步。如果发现 160ms 花费在外部积分服务调用上那么优化思路就应该是把同步调用改成异步消息、增加超时控制、或者做批量合并调用。4. 更成熟的优化方案与代码示例4.1 版本说明本文示例基于常见环境编写主要演示思路和关键配置。示例使用 Kafka 2.x/3.x 客户端 APISpring Boot 项目基于 2.x/3.x 均可运行部分参数在不同版本中名称略有差异请以你项目实际使用的版本为准。建议环境如下JDK 8 或 11。Spring Boot 2.7 或 3.x。Kafka 客户端版本 2.8。Maven 3.6。4.2 方案一消费逻辑瘦身这是最优先做的优化成本最低效果往往最明显。Kafka 消费者在拉取一批消息后会逐条执行业务逻辑。如果业务逻辑中包含以下操作会明显拖慢消费速度每条消息都执行一次数据库写操作。每条消息都调用一次外部 HTTP 接口。每条消息都发送一条短信或推送通知。在消费线程中执行了耗时很长的计算任务。针对这些情况可以做三件事第一尽可能把非核心操作异步化。比如发短信、推送通知完全可以在消费逻辑中先落库再投递到另一个 topic由专门的消费者去处理。第二减少外部调用次数。可以把多条消息合并成一批批量调用外部接口或者把一条消息中的多个外部依赖并行调用而不是串行调用。第三过滤无效消息。很多场景下队列中存在大量无需处理的消息比如重复通知、测试消息、过期数据。在消费入口处增加过滤逻辑可以直接降低无效消费。下面是一个简化版的消费监听器演示如何快速过滤无效消息// 文件路径src/main/java/com/example/kafka/OrderConsumer.java Component public class OrderConsumer { private static final Logger log LoggerFactory.getLogger(OrderConsumer.class); KafkaListener(topics order-topic, groupId order-consumer-group) public void onMessage(ConsumerRecordString, String record) { OrderEvent event JSON.parseObject(record.value(), OrderEvent.class); if (event null || event.getOrderId() null) { // 无效消息直接记录并跳过 log.warn(invalid message, offset{}, record.offset()); return; } if (event.getEventTime() System.currentTimeMillis() - 30 * 60 * 1000L) { // 超过30分钟的消息业务上已无处理意义跳过 log.info(expired message ignored, orderId{}, event.getOrderId()); return; } // 真正的业务处理 process(event); } private void process(OrderEvent event) { // 模拟业务处理 } }这里的核心思想是消费入口要先做轻量级判断把不需要处理的消息挡在门外避免无意义的耗时。4.3 方案二批量消费与手动提交很多团队默认使用 Spring Boot 的KafkaListener逐条消费每条消息处理完自动提交 offset。这种方式实现简单但性能有限。原因有两个第一逐条拉取和逐条提交会增加网络与 broker 的交互次数。 第二自动提交模式下消费者每处理一条就提交一次如果处理失败很容易造成 offset 频繁提交和回滚。更高效的做法是开启批量消费并改为手动提交 offset。这样每次 poll 可以拉取一批消息处理完成后统一提交一次 offset大幅减少提交次数。在 Spring Boot 中可以通过application.yml做如下配置spring: kafka: bootstrap-servers: localhost:9092 consumer: group-id: order-consumer-group enable-auto-commit: false auto-offset-reset: latest key-deserializer: org.apache.kafka.common.serialization.StringDeserializer value-deserializer: org.apache.kafka.common.serialization.StringDeserializer max-poll-records: 500 listener: type: batch ack-mode: manual_immediate concurrency: 3配置说明enable-auto-commit: false关闭自动提交避免消息未处理完就提交 offset。max-poll-records: 500单次 poll 最多拉取 500 条消息减少网络往返。listener.type: batch监听器按批量模式处理消息。ack-mode: manual_immediate手动调用 ack 时立即提交保证灵活性。concurrency: 3并发消费者线程数建议不要超过分区数。批量消费者代码示例如下// 文件路径src/main/java/com/example/kafka/BatchOrderConsumer.java Component public class BatchOrderConsumer { private static final Logger log LoggerFactory.getLogger(BatchOrderConsumer.class); KafkaListener(topics order-topic, groupId order-consumer-group) public void onBatchMessage(ListConsumerRecordString, String records, Acknowledgment ack) { long start System.currentTimeMillis(); try { ListOrderEvent events new ArrayList(records.size()); for (ConsumerRecordString, String record : records) { OrderEvent event JSON.parseObject(record.value(), OrderEvent.class); if (event ! null event.getOrderId() ! null) { events.add(event); } } // 批量业务处理例如批量写库 orderService.batchProcess(events); // 处理成功手动提交 offset ack.acknowledge(); log.info(batch consume success, size{}, cost{}ms, records.size(), System.currentTimeMillis() - start); } catch (Exception e) { // 处理失败不建议无限重试先记录并提交避免阻塞后续消息 log.error(batch consume error, size{}, records.size(), e); ack.acknowledge(); } } }这里需要注意手动提交 offset 的时机。如果业务处理失败后直接提交可能会有消息丢失风险如果不提交又会导致该批消息反复拉取消费进度无法前进。在生产环境中更稳妥的做法是先对消息做幂等处理处理失败的消息投递到死信 topic当前批次正常提交。4.4 方案三线程池并发消费模型批量消费能减少网络交互次数但如果单条消息的核心业务逻辑很重比如每个订单都要调用积分服务那即便批量拉取最终还是逐条串行处理吞吐提升有限。这时可以引入线程池并发消费模型消费者线程只负责拉取消息和提交 offset真正的业务处理交给独立的业务线程池执行。一个简单的并发消费模型如下// 文件路径src/main/java/com/example/kafka/ConcurrentOrderConsumer.java Component public class ConcurrentOrderConsumer { private static final Logger log LoggerFactory.getLogger(ConcurrentOrderConsumer.class); private final ExecutorService bizExecutor new ThreadPoolExecutor( 8, 16, 60L, TimeUnit.SECONDS, new ArrayBlockingQueue(2000), new ThreadPoolExecutor.CallerRunsPolicy() ); KafkaListener(topics order-topic, groupId order-consumer-group) public void onMessage(ConsumerRecordString, String record, Acknowledgment ack) { bizExecutor.execute(() - { try { OrderEvent event JSON.parseObject(record.value(), OrderEvent.class); // 真正的业务处理可以在这里做批量合并、异步化等操作 orderService.process(event); } catch (Exception e) { log.error(biz process error, offset{}, record.offset(), e); } }); // 注意这里提交 offset 的前提是业务已做幂等允许重复消费 ack.acknowledge(); } }这个模型的核心收益是poll 线程不会被业务逻辑阻塞可以持续从 Kafka 拉取新消息从而提升整体消费速率。但必须强调这种模型有三个重要前提第一业务逻辑必须幂等。因为消费者拉取到消息后可能还没有处理完offset 就已经提交了。如果消费者宕机这部分消息会被重新拉取导致重复处理。第二线程池队列不能无限增长。如果下游处理速度跟不上线程池队列会积压大量任务最终导致内存溢出。建议使用有界队列并设置合理的拒绝策略。第三顺序性需要额外设计。如果业务要求同一订单的消息必须严格按照顺序处理那么简单使用多线程并发消费会破坏顺序。这时可以按订单 ID 进行哈希将同一订单的所有消息路由到同一个线程处理。4.5 方案四同步改异步下游削峰在很多积压场景中真正的瓶颈不是 Kafka 本身而是下游系统。比如消费一条订单消息需要执行以下操作写订单表。写流水表。扣减库存。调用积分服务。发送站内通知。其中写订单和写流水是核心操作必须尽快完成。而调用积分服务、发送通知属于非核心操作完全可以在消费主链路中砍掉。改造思路消费主流程只做必要的数据落库。把非核心操作发送到另一个 topic由专门的消费者异步处理。对必须调用的外部接口增加超时控制和降级逻辑。数据库写入尽量使用批量 insert减少事务开销。批量写库示例// 文件路径src/main/java/com/example/kafka/OrderService.java Service public class OrderService { Autowired private OrderMapper orderMapper; Transactional(rollbackFor Exception.class) public void batchProcess(ListOrderEvent events) { if (events null || events.isEmpty()) { return; } // 批量插入减少单条 insert 的事务开销 ListOrderDO orders events.stream() .map(this::toDO) .collect(Collectors.toList()); orderMapper.batchInsert(orders); // 其他核心操作... } }如果批量插入的条数太多建议分批执行比如每 500 条作为一个事务避免单个事务过大导致锁竞争。4.6 方案五动态感知消费压力临时降速或拒绝还有一种场景是消费者本身处理能力没问题但上游生产者短时间内发送了过多消息形成瞬时洪峰。比如大促期间订单量突增消费速度跟不上生产速度。这种情况下可以通过以下方式缓解在消费入口增加速率控制例如使用 Guava RateLimiter 限制每秒最大处理条数。如果消息时效性要求不高可以临时提高max.poll.records让单次拉取更多消息减少 poll 次数。如果下游已经处于过载状态可以考虑将部分消息投递到备份 topic等高峰期过去后再消费。这里需要说明的是速率控制本质上是在“损失消费速度”和“保护下游”之间做权衡不能作为长期方案。长期来看还是需要提升消费能力或者优化业务链路。5. 实战一条消息 200ms如何把积压追平速度提升 10 倍5.1 场景与瓶颈确认假设线上有一个订单消息 topic共 20 个分区消费者组名为order-consumer-group部署了 5 个消费者实例。高峰期每分钟产生 6000 条订单消息平均每秒 100 条。业务方反馈消息处理延迟越来越大Lag 最高达到 80 万条。排查时先查看消费组详情kafka-consumer-groups.sh \ --bootstrap-server localhost:9092 \ --group order-consumer-group \ --describe输出显示每个分区 LAG 都很高且持续增长。再看消费日志发现单条消息平均处理耗时 200ms其中解析消息5ms。查询数据库关联数据20ms。调用积分服务150ms。写库20ms。其他5ms。可以算出单个消费者单线程每秒只能处理 5 条消息。5 个消费者对应 20 个分区即使每个消费者都处理 4 个分区也没有能力并行处理一个分区内的多条消息。因为单个分区内的消息是串行消费的所以单分区每秒只能处理 5 条。整个消费组理论上限是 100 条/秒但受制于积分服务耗时实际只有 25 条/秒左右远低于生产速度 100 条/秒。注意这里要额外说明一点单分区内消息虽然是顺序拉取但如果消费者使用线程池并行处理同一分区的消息也可以并行处理只是会牺牲顺序性。本例中我们采用批量消费 异步化方案。5.2 改造前代码改造前的消费者监听器大致如下// 文件路径src/main/java/com/example/kafka/OldOrderConsumer.java Component public class OldOrderConsumer { private static final Logger log LoggerFactory.getLogger(OldOrderConsumer.class); KafkaListener(topics order-topic, groupId order-consumer-group) public void onMessage(ConsumerRecordString, String record) { long start System.currentTimeMillis(); OrderEvent event JSON.parseObject(record.value(), OrderEvent.class); OrderDO order convert(event); // 查询关联数据 UserDO user userMapper.selectById(event.getUserId()); // 调用积分服务 int points pointsService.addPoints(user.getId(), event.getAmount()); // 写库 orderMapper.insert(order); // 发通知 notifyService.send(event.getUserId(), 订单处理成功); log.info(consume one message, orderId{}, cost{}ms, event.getOrderId(), System.currentTimeMillis() - start); } }这段代码中耗时最大的就是pointsService.addPoints这个同步调用耗时约 150ms。5.3 改造后代码与配置改造分两步。第一步把积分服务调用和通知发送从主链路中拆出去。消费主流程只保留关联数据查询和订单落库。第二步开启批量消费数据库写入改为批量 insert。改造后的消费者// 文件路径src/main/java/com/example/kafka/NewOrderConsumer.java Component public class NewOrderConsumer { private static final Logger log LoggerFactory.getLogger(NewOrderConsumer.class); KafkaListener(topics order-topic, groupId order-consumer-group) public void onBatchMessage(ListConsumerRecordString, String records, Acknowledgment ack) { long start System.currentTimeMillis(); ListOrderEvent events new ArrayList(records.size()); for (ConsumerRecordString, String record : records) { OrderEvent event JSON.parseObject(record.value(), OrderEvent.class); if (event ! null event.getOrderId() ! null) { events.add(event); } } if (events.isEmpty()) { ack.acknowledge(); return; } // 主流程批量落库 orderService.batchInsertOrders(events); // 异步发送积分服务和通知消息 for (OrderEvent event : events) { kafkaTemplate.send(order-points-topic, event.getOrderId(), event); kafkaTemplate.send(order-notify-topic, event.getOrderId(), event); } ack.acknowledge(); log.info(batch consume success, size{}, cost{}ms, records.size(), System.currentTimeMillis() - start); } }对应地在application.yml中增加批量配置spring: kafka: bootstrap-servers: localhost:9092 consumer: group-id: order-consumer-group enable-auto-commit: false auto-offset-reset: latest max-poll-records: 500 key-deserializer: org.apache.kafka.common.serialization.StringDeserializer value-deserializer: org.apache.kafka.common.serialization.StringDeserializer listener: type: batch ack-mode: manual_immediate concurrency: 55.4 运行验证与效果说明改造后原本 200ms 的单条处理链路被拆分为批量拉取 500 条消息。批量解析与批量写库500 条总耗时约 1.2 秒。积分服务和通知通过异步 topic 继续处理。简单估算一下改造前单消费者吞吐 1000ms / 200ms 5 条/秒 消费组总吞吐 5 条/秒 × 5 个消费者 25 条/秒受单分区串行限制实际可能更低改造后单批处理耗时 1200ms / 500 条 2.4ms/条 单消费者吞吐 ≈ 400 条/秒 消费组总吞吐 ≈ 400 条/秒 × 5 个消费者 2000 条/秒如果分区和网络允许虽然这个数值是理想估算实际还要考虑数据库写入能力和网络开销但吞吐量提升 10 倍以上是完全可以做到的。80 万条积压消息在 2000 条/秒的消费速度下大约 400 秒就能追平也就是 7 分钟左右。相比改造前需要 8 小时以上差距非常明显。需要注意的是这个方案里order-points-topic和order-notify-topic的消费能力也要纳入整体规划。如果这两个 topic 的消费者也存在瓶颈需要同样优化或者接受“主流程快、副流程慢”的状态。6. 什么时候扩容仍然是正确选择前面说了很多扩容的局限性但并不意味着扩容在任何场景下都不应该用。成熟工程师的做法是“按需扩容、算清楚再扩”。下面几种情况下扩容是合理且必要的。6.1 分区数成为硬性上限根据 Kafka 的分区消费模型一个分区同一时间只能被消费者组内的一个消费者实例消费。假设某个 topic 只有 3 个分区那么即使你部署 10 个消费者实例实际并行消费的也只有 3 个。这种情况下如果消费者实例的并发度已经达到上限且消费逻辑已经优化过瓶颈仍然存在那么合理的做法是增加分区数。增加分区数的命令kafka-topics.sh \ --bootstrap-server localhost:9092 \ --alter \ --topic order-topic \ --partitions 40但是增加分区数有几个重要注意事项第一分区数只能增加不能减少。虽然 Kafka 3.x 支持减少分区但操作复杂且有很多限制生产环境不要轻易尝试。第二增加分区会改变消息分布。如果生产者使用 key 来保证某个 key 的消息进入同一分区那么增加分区后key 与分区的映射关系会变化可能影响消息顺序。第三增加分区后需要确认消费者实例数没有超过分区数否则部分消费者实例会空闲。因此扩容分区数的前提是确认业务可以接受分区重新分布带来的顺序性变化且消息量确实需要更多分区来承载并发。6.2 消费者节点资源确实不足如果消费者节点的 CPU、内存、网络带宽已经达到瓶颈且无法通过优化消费逻辑来解决那么扩容消费者实例是有效的。比如消费者节点 CPU 长期 90% 以上说明单机处理能力已经饱和。此时增加消费者实例把分区分摊到更多机器上能明显提升总吞吐。扩容时建议小步快跑每次增加 1 到 2 个实例观察 Lag 和 rebalance 情况再决定是否继续增加。6.3 扩容的正确姿势即使决定扩容也不是简单加机器就完事。建议按以下顺序操作确认分区数是否足够。如果分区数不足先增加分区数。检查消费者组内是否有空闲消费者。如果有说明当前消费配置不合理先调整 concurrency。小步增加消费者实例观察 rebalance 频率。如果 rebalance 次数明显增加说明新增实例可能引发了分区频繁转移。扩容后持续观察 Lag 指标确认积压在下降而不是短暂回升后停止。如果扩容后 Lag 没有明显下降及时回滚变更重新排查瓶颈。7. 生产环境最佳实践与高频问题排查7.1 积压治理中的工程规范结合上述分析在真实项目中治理 Kafka 积压问题下面这些规范很值得沉淀到团队里。第一消费逻辑必须幂等。Kafka 只能保证“至少一次”投递不能保证“恰好一次”。消费者在异常重启、rebalance、手动提交 offset 时都可能收到重复消息。幂等设计是消费端最基本的要求。常见做法是使用唯一业务键查重或者建一张去重表。第二不要关闭自动提交就直接上线。从自动提交改成手动提交时要考虑消息处理失败后的策略。如果失败就无限重试会导致消费者卡死Lag 持续上涨。更合理的做法是设置重试上限超过上限进入死信 topic。第三监控指标要完善。至少监控以下指标消费组 Lag、消费速率、单条消息处理耗时、rebalance 次数、消费者线程活跃数、下游系统耗时。建议通过日志或 Micrometer 主动上报这些指标。第四消费线程与业务线程分离。不要让耗时的业务逻辑直接阻塞 Kafka 的 poll 线程。使用独立的业务线程池时要设置合理的队列大小与拒绝策略。第五生产变更必须小流量验证。无论是修改消费并发数、调整max.poll.records还是增加分区都应该先在测试环境验证再逐步灰度到生产。变更前做好回滚预案。第六数据库批量操作要控制粒度。批量 insert 不是越大越好建议每批 200 到 500 条单事务执行时间控制在秒级以内避免长时间占用数据库连接和锁资源。7.2 高频问题排查表下面是 Kafka 消息积压场景中比较常见的几个问题现象和排查方向。问题现象常见原因解决思路Lag 持续增长消费速率很低单条消息处理耗时长业务逻辑有同步外部调用优化消费逻辑、异步化、减少外部调用Lag 增长但日志中出现大量 rebalance消费者处理超时触发 max.poll.interval.ms调整消费超时配置把业务逻辑移出 poll 线程增加消费者实例后 Lag 没有下降消费者实例数已经超过分区数或分区数不足增加分区数或调整 concurrency某个分区 Lag 特别大其他分区正常分区数据倾斜或该分区消费者实例异常检查 key 设计确认是否有热点 key检查对应消费者日志消费者频繁报错消费不断重试下游接口不稳定或反序列化失败增加重试上限失败消息投递死信 topicCPU 使用率不高但消费速率慢瓶颈在下游数据库或外部接口使用链路追踪定位耗时优先优化下游批量消费后消息丢失业务处理失败但 offset 已提交失败消息必须进入重试或死信队列不能静默丢弃修改配置后复现积压配置参数与业务不匹配如 max.poll.records 过大按业务处理能力设置合理参数单批处理时间控制在 max.poll.interval.ms 内7.3 一个可以直接使用的检查清单如果你今天接到一个“Kafka 积压”的线上告警可以按下面的顺序排查先看消费组 Lag 分布确认是整体积压还是单分区积压。看消费者日志是否有 rebalance 或异常重试。如果有链路追踪看单条消息耗时分布定位耗时瓶颈。如果瓶颈是外部调用考虑异步化或批量合并。如果瓶颈是数据库写入考虑批量写、合并写、分库分表。如果所有优化都做过且资源确实不足再考虑扩容。扩容前确认分区数扩容后观察 rebalance 和 Lag 变化。处理完积压后复盘根因并完善监控告警。8. 总结与行动建议Kafka 积压问题的本质是生产速度与消费速度不匹配但匹配失衡的原因可能出现在消息生产、消费逻辑、下游依赖、资源配置等多个环节。初学者遇到积压第一反应是扩容而有经验的工程师会先做指标分析、耗时拆解和瓶颈定位再做针对性优化。这篇文章的核心结论可以浓缩成几句话扩容只对“消费者自身处理能力不足”这一种原因有效。消费者数量超过分区数时扩容无效。优化消费逻辑、开启批量消费、使用线程池并发消费、同步改异步通常是更优先的解决手段。分区数扩容要谨慎评估顺序性影响且分区数通常只能增加不能减少。生产环境必须配合指标监控、幂等消费和死信队列机制才能长期避免积压问题。如果你下次再遇到 Kafka 积压告警建议先冷静 10 分钟查一下消费组 Lag 分布看一下单条消息耗时再决定到底是要改代码、调配置还是加机器。多数情况下你会发现代码优化比扩容更有效成本也更低。把排查步骤和优化方案整理成团队文档下次再遇到类似问题的时候就能快速定位不用每次都在线上一通乱试。