十年匠心定制 · 商业建站与技术教学双线并行 咨询热线:400-886-1026 service@lmnt.cn
ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

Kafka积压别急着扩容:先定位这三步再动手

Kafka积压别急着扩容:先定位这三步再动手 线上 Kafka 积压告警你打开电脑登录控制台准备给消费者组扩容——如果你是这么做的那这篇文章就是写给你的。先说结论遇到 Kafka 积压就扩容是典型的初学者思维。扩容不是不能解决积压而是它往往解决的不是积压的根因只是把问题往后推了几个小时。真正有价值的处理方式是先回答三个问题消息积压在哪个环节消费者慢是慢在哪一次调用上这个消费者组的并行度天花板在哪里这篇文章不会只讲概念。我会从 Kafka 的消费模型讲起然后给你一套定位积压的排查路径再按不同的积压成因给出对应的处理方案最后用命令行、配置和代码演示如何落地。读完你可以直接照着排查一次自己项目的积压问题。1. 为什么会积压先理解 Kafka 的消费模型很多人对 Kafka 积压的理解停留在消息太多了消费不过来这个说法没有错但它掩盖了真正的技术细节。要搞清楚积压必须先理解 Kafka 消费的三个核心概念分区、消费者组、offset。分区Partition是 Kafka 并行度的基础。一个 Topic 被拆成多个分区每个分区内的消息是有序的。生产者在写消息时可以指定 key相同 key 的消息会进入同一个分区目的是保证局部有序。消费者组Consumer Group是 Kafka 最巧妙的设计之一。同一个消费者组内的多个消费者实例协同消费同一个 Topic组内的分区会被分配到这个组的不同消费者上。这里有一个很多人没意识到的硬性规则一个分区在同一时刻只能被同一个消费者组内的一个消费者实例消费。这意味着什么假设你有一个 Topic 只有 3 个分区消费者组里有 10 个消费者实例那么同一时刻最多只有 3 个消费者在干活另外 7 个消费者完全是空闲的。你扩容到 20 个也一样3 个在干活17 个在空闲。这就是扩容无效的第一个典型场景——分区数决定了消费者的并行度上限。offset偏移量是消费者在某个分区上的消费位置标记。消费者每次拉取一批消息处理完成后提交 offset下次拉取就从新的 offset 开始。积压的本质就是最新写入的消息位置Log End OffsetLEO和消费者当前消费位置Current Offset之间的距离在持续拉大。这个距离通常被称为 Lag。Lag 持续增长说明消费速度长期低于生产速度。理解了这三个概念你就会明白解决积压的本质要么是提升消费速度要么是降低单位时间内的消息量要么是扩大并行度。扩容只是扩大并行度的一种手段而且是很有限的手段。2. 导致积压的几种真实原因从实际项目排查经验看Kafka 积压的原因可以归纳为四类。每一类背后对应的解法完全不同。2.1 消费逻辑本身太慢这是最常见的原因也是最容易被甩锅给Kafka 慢的原因。消费者把消息拉下来之后处理逻辑是一个重量级操作写数据库、调用第三方接口、批量计算、发通知。单条消息处理耗时从几十毫秒涨到几百毫秒消费速度自然跟不上生产速度。这里要区分一个概念消费慢不等于 Kafka 慢。消息从 Broker 拉取到消费者的过程通常是毫秒级真正耗时的是你业务代码中的那一堆逻辑。遇到这种情况你来回来去调 Kafka 参数是没有意义的应该把注意力放到业务代码性能分析上。2.2 消费者并行度或分区数不足如果你的消费逻辑本身不慢但一个消费者组只有一个消费者实例而这个组订阅的 Topic 有 50 个分区那么这 50 个分区的所有消息都由一个消费者处理。单个消费者的处理能力是有上限的无论单条消息处理多快都容易积压。还有一种情况是消费者实例数量多于分区数。比如 10 个消费者实例但 Topic 只有 6 个分区那么 4 个消费者长期空闲。你以为自己在并行消费实际上只有 6 个消费者在干活。这种情况你看到积压后再增加消费者实例数量不会有任何效果因为分区数的上限已经卡住了。2.3 消费者处理出现异常频繁重试或直接卡死积压只是表象背后的真相可能是消费者代码抛异常了。很多团队在 catch 块里做重试重试失败写日志但 offset 没有提交下次重启又从头消费同一批数据。如果异常总是发生消费就会在原地打转Lag 一路上涨。更隐蔽的一类是消费者线程阻塞。比如消费者代码里有一个同步调用超时时间设置得太长或者数据库连接池满了消费者线程全部阻塞在等待上。此时消费者既没有崩溃也没有明显报错但就是一动不动。2.4 Broker 侧问题比如某个分区所在磁盘接近写满读写性能下降或者某个 Broker 节点负载过高或者副本同步出现延迟导致 ISR 列表收缩leader 切换频繁。这些情况也会影响消费速度但相对少见而且通常会伴随其他监控指标异常。排查积压时不能只盯着消费者看也要确认 Broker 侧是否有异常。积压原因分类典型表现排查方向解决思路消费逻辑慢Lag 平稳增长CPU/IO无瓶颈打印单条消息处理耗时优化业务代码、批量处理分区数不足消费者数大于分区数部分消费者空闲查看分区分配情况增加分区数、调整消费者数量异常重试卡死offset 不更新日志反复报错查看消费者异常日志修复异常处理逻辑、配置重试策略Broker 侧问题所有消费者都变慢集群指标异常查看 Broker 负载和磁盘调整分区分布、优化集群配置3. 先定位如何判断积压卡在哪一环处理积压的第一步不是动手解决而是定位。推荐按这个顺序排查先确认为什么消费速度变慢再确认是单消费者卡住还是整体变慢最后确认是消费逻辑还是 Broker 侧问题。3.1 用命令行查看消费者组 LagKafka 自带的命令行工具是排查积压的首选工具。不需要额外部署也不需要写代码直接在 Kafka 服务器上执行即可。# 查看指定消费者组的消费情况和积压量 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-consumer-group test-topic 0 1024 2048 1024 consumer-1-xxx your-consumer-group test-topic 1 1024 2048 1024 consumer-1-xxx your-consumer-group test-topic 2 2000 2048 48 consumer-2-xxx这里你要看的不是所有分区的 Lag 是不是大于 0而是看 Lag 的分布。所有分区的 Lag 都很大可能是消费逻辑整体变慢或者生产者消息量突增或者某个外部依赖变慢影响了所有分区的消费。部分分区 Lag 很大部分正常这通常说明分配在某些分区上的消费者实例处理较慢或者某些分区的数据存在热点。某个分区的 CURRENT-OFFSET 长时间不动消费者可能在这个分区上卡住了。如果你的消费代码是同步处理模型一个分区的消息会阻塞后续同分区消息的消费。这个命令只能看到某一时刻的快照。要判断 Lag 是在增长还是在收敛需要隔几分钟跑一次对比或者直接接入监控系统。3.2 确认消费者实例数量和分布接着确认消费者组里的实际实例数以及每个实例分配到了哪些分区。# 详细查看消费者实例和分区分配 kafka-consumer-groups.sh --bootstrap-server localhost:9092 --group your-consumer-group --describe --members这个命令很关键。你会清楚地看到组内有多少个活跃消费者实例每个实例分配到了哪些分区。如果 Topic 有 10 个分区而组内只有 2 个消费者实例那每个消费者要承担 5 个分区的消费任务。此时你就要评估单消费者实例的消费能力是否能跟上。3.3 给消费链路加上耗时日志命令行能看出慢但不能直接告诉你为什么慢。最直接的办法是在消费者代码里加一段简单的耗时统计把消息从拉取到完成处理的每个阶段耗时打出来。// 伪代码关键思想是分段计时 long start System.currentTimeMillis(); // 阶段1拉取后的预处理 ListRecord records pollRecords(); long afterPoll System.currentTimeMillis(); log.info(poll {} records, cost {} ms, records.size(), afterPoll - start); // 阶段2执行业务逻辑 for (Record record : records) { process(record); } long afterProcess System.currentTimeMillis(); log.info(process cost {} ms, afterProcess - afterPoll); // 阶段3提交offset consumer.commitSync(); long afterCommit System.currentTimeMillis(); log.info(commit cost {} ms, afterCommit - afterProcess);如果process cost动辄几百毫秒那答案已经很明确了业务逻辑跟不上。如果poll本身耗时很高那可能是网络或者 Broker 侧问题。这种最低成本的分段埋点往往比上全套监控系统更快定位问题。4. 为什么扩容不是首选什么情况才真的需要扩容现在可以正面回答标题里的问题了。4.1 扩容机器的三个副作用如果你在没有定位根因的情况下直接扩容消费者实例通常会遇到以下情况。第一新增的消费者实例根本不干活。分区数少于消费者实例数时新增的消费者不会分配到任何分区属于纯资源浪费。扩容前后Lag 指标不会有任何变化但你的机器成本和维护成本实实在在增加了。第二扩容短期掩盖了问题但根因没解决。假设消费逻辑单条消息要 500ms你把消费者从 2 台扩到 10 台。如果分区数足够多消费速度确实提升了Lag 下降了。但是如果那 500ms 是因为调用外部接口超时那么扩容之后对外部接口的并发调用量同步增加很可能直接把外部接口打挂。这个后果比积压本身严重得多。第三状态与运维复杂度增加。消费者实例变多意味着你需要在配置中心、监控系统、日志系统里同步增加节点信息。一旦消费者代码有 bug排查范围也从 2 台变成 10 台。4.2 真正适合扩容的场景说这些不是完全否定扩容。在下面几种场景中扩容是合理且必要的。场景一分区数量足够消费者实例数跟不上。你确认 Topic 有 30 个分区但消费者只有 2 个实例。此时增加消费者实例到 5 个甚至 10 个可以让并行度成倍提升这是扩容最直接有效的场景。场景二短期内消息量突增比如大促、秒杀、数据补偿任务。业务方在某个时段会集中推送大量消息正常存量消费能力确实无法覆盖。这种情况下扩容是应急手段但要同步确认下游依赖能承受住增加的压力并且要在流量过去后及时缩容。场景三某个消费者实例所在节点出现资源问题。比如 CPU、内存、磁盘 IO 已经打满此时在原有节点上硬扛没有意义加一台机器分流是合理的。4.3 扩容的正确姿势检查分区数和实例数是否匹配在动手扩机器之前先用下面的思路确认扩容是否有效。1. 查看 Topic 分区总数。 2. 查看消费者组内活跃实例数。 3. 如果 分区数 实例数扩容有效。 4. 如果 分区数 实例数扩容无效先调整分区或定位消费逻辑。如果分区数和实例数已经接近甚至超过那就不是扩容的问题而是代码或架构层面需要调整。5. 消费端代码级调优比扩容更优先做的事在实际项目中消费端代码优化通常比扩容收益更大。下面用 Spring Boot 集成 Kafka 的场景做一个完整示例。5.1 检查并调整消费者关键参数如果你的消费者应用是 Spring Boot 项目最常用的配置是application.yml中的 Spring Kafka 消费者配置。spring: kafka: bootstrap-servers: localhost:9092 consumer: group-id: order-consume-group enable-auto-commit: false key-deserializer: org.apache.kafka.common.serialization.StringDeserializer value-deserializer: org.apache.kafka.common.serialization.StringDeserializer max-poll-records: 500 fetch-max-wait-ms: 500 fetch-min-bytes: 1048576 session-timeout-ms: 15000 max-partition-fetch-bytes: 1048576 auto-offset-reset: latest listener: concurrency: 3 type: batch ack-mode: manual-ack先解释几个关键配置。max-poll-records单次poll最多拉取多少条消息。拉得太多业务处理不过来反而会超过max.poll.interval.ms导致消费者被判定离线拉得太少每次网络往返的收益太低消费速度上不去。从经验看如果每条消息处理较快500 是一个合理起点。enable-auto-commit: false和ack-mode: manual-ack关闭自动提交改为手动提交。很多积压和重复消费问题都源于自动提交导致的 offset 和业务处理不同步。listener.concurrencySpring Kafka 创建多少个消费者线程。注意这个值不是想设多大就多大分区总数是硬约束。如果分区数是 6concurrency设为 10最终只有 6 个线程能分配到分区。5.2 一个正确处理手动提交的消费者代码手动提交的核心原则是消息先处理成功再提交 offset。反过来如果先提交 offset处理失败时消息就丢了。// 文件路径src/main/java/com/example/kafka/OrderConsumer.java Slf4j Component public class OrderConsumer { KafkaListener(topics order-topic, groupId order-consume-group) public void onMessage(ListConsumerRecordString, String records, Acknowledgment ack) { ListOrder orders new ArrayList(records.size()); log.info(poll message size: {}, records.size()); // 1. 解析消息 long start System.currentTimeMillis(); for (ConsumerRecordString, String record : records) { Order order JSON.parseObject(record.value(), Order.class); orders.add(order); } log.info(parse order cost {} ms, System.currentTimeMillis() - start); // 2. 批量写入数据库 long dbStart System.currentTimeMillis(); if (!orders.isEmpty()) { orderService.batchInsert(orders); } log.info(batch insert cost {} ms, System.currentTimeMillis() - dbStart); // 3. 确认消息提交offset ack.acknowledge(); } }这段代码有几个细节值得注意。批量处理而非单条处理。如果每条消息单独写库500 条消息就要发起 500 次网络请求。改成批量写库耗时可能只有原来的 1/5 到 1/10。这是处理积压最有效的手段之一。正确处理 ack。只有当批量处理全部成功之后才调用ack.acknowledge()。如果中间出现异常可以不 ack让消费者重新拉取这批消息同时配合错误日志告警。5.3 调优后的性能提升预期这里不是拍脑袋而是给一个判断逻辑。假设处理一条消息需要 10ms单线程每秒可以处理 100 条。如果把并发数从 1 提升到 3消费速度会接近每秒 300 条。如果进一步从单条处理改成批量处理批量 100 条写一次库单次耗时可能只需要 50ms相当于每条消息摊 0.5ms消费速度提升接近 20 倍。这就是先优化代码再考虑扩容的意义所在。6. 生产者侧和 Broker 侧导致的积压处理不是所有积压都发生在消费端。如果你确认消费逻辑没问题并发也够但 Lag 还是增长那就要检查生产者和 Broker 侧。6.1 生产端热 key 导致的分区倾斜当一个 Topic 的某个 key 消息量特别大而 Kafka 的 key 哈希策略会把相同 key 的消息路由到同一个分区就会出现分区倾斜某一个分区的数据量远大于其他分区。此时你用命令行看 Lag会看到部分分区很大部分分区正常。处理思路一般是修改 key 设计。比如原本用用户 ID 作为 key 保证单个用户消息有序现在发现某个大用户的订单量过大可以把 key 改成用户 ID 分片后缀。这样既能在大多数场景下保持同一用户的局部有序又避免了单个分区过载。6.2 单分区顺序消费的瓶颈有些业务出于顺序要求必须保证同一个 key 的消息经过同一个分区被顺序消费。如果这种 key 的粒度很粗比如整个订单系统的消息都用orderId但都落到少数几个分区上顺序是保住了但消费吞吐量被严重限制。解决思路通常需要业务上做权衡是否真的需要全局有序如果可以容忍短时间乱序就可以放宽 key 粒度。6.3 Broker 分区副本分布不均用 Kafka 自带的命令可以查看分区的副本分布情况。kafka-topics.sh --bootstrap-server localhost:9092 --describe --topic order-topic如果大量分区的 leader 集中在同一个 Broker 上那么这个 Broker 的读写压力会非常大。此时可以用kafka-reassign-partitions.sh做分区均衡把部分分区的 leader 分散到其他 Broker 上。这个是集群运维层面的操作执行前要确认当前集群的负载情况并且建议先在测试环境验证。7. 验证效果确认积压正在恢复而不是暂时掩盖解决积压方案部署上线之后不要急着收工。你需要用数据确认积压确实在收敛而不是换了种方式继续积压。7.1 多次执行 Lag 查询观察趋势这是最直接的验证方式。在方案上线 5 分钟、15 分钟、30 分钟、1 小时各执行一次前面的kafka-consumer-groups.sh命令记录每个分区的 LAG。LAG 持续下降说明消费速度正在超过生产速度方案有效。LAG 先降后涨说明生产端流量可能持续上涨或者消费端仍然存在瓶颈。LAG 完全不动说明新增的消费者实例很可能没有分配到分区或者代码里出现了持续性异常。7.2 查看消费者日志中的持久化耗时如果你在代码里加了耗时日志这是最能直接说明问题的数据。当 Lag 开始下降时日志里的process cost应该保持在合理范围比如单批消费耗时不超过 1 秒。如果 Lag 下降但process cost还在持续升高说明你的消费代码可能引入了新的阻塞点需要继续排查。7.3 用告警而不是人为巡检真正可靠的方式是配置消费 Lag 监控和告警。无论用 Kafka 自带的 JMX 指标还是接入 Prometheus、Grafana、Kafka Eagle 等可视化工具只要消费者的 Lag 超过阈值就自动触发告警。人为巡检只适合应急场景生产环境必须靠自动化监控。这里有一个容易忽略的点不要在积压发生时才看监控应该在平时的消费延迟基线之上设置告警阈值。比如正常情况下 Lag 长期在几百以内那阈值可以设到 5000如果阈值设得过低会因为正常抖动频繁误报时间久了团队就会麻木。8. 常见问题与排查思路下面把处理 Kafka 积压过程中经常遇到的问题整理成一个排查表方便收藏备用。问题现象可能原因排查方式解决方案所有分区 Lag 同步增长消息量突增或业务逻辑整体变慢对比生产速率、消费耗时日志优化消费逻辑、批量处理、必要时扩容部分分区 Lag 很大部分正常分区数据倾斜或消费实例分配不均查看分区详情和成员分布调整 key 设计、重新分配分区新增消费者实例但没有生效消费者实例数超过分区数--describe --members查看分配增加分区数或减少消费者实例LAG 不降但消费者无异常日志消费者线程阻塞在外部调用查看线程 dump、接口超时时间优化超时配置、增加线程池隔离offset 提交后消息重复处理处理逻辑不在 ack 之前完成检查代码中 ack 顺序先处理成功再 ack消费者反复加入退出rebalance 频繁单个消费者处理耗时过长超过 max.poll.interval.ms查看 rebalance 日志和单批处理耗时降低 max.poll.records、异步化处理、延长超时某个 Broker 磁盘或 CPU 过高分区分布不均或 Broker 资源问题查看集群监控和分区分布分区均衡、磁盘扩容、迁移分区9. 最佳实践与工程建议最后分享几条 Kafka 消费者工程落地时非常有价值的实践经验。这些经验看起来简单但每一条背后都对应过真实的线上事故。第一topic 创建时就要规划好分区数。分区数一旦确定后续调整成本很高。建议消息量较小的 Topic 分区数设 6 到 12 个消息量明显较大的 Topic 设 24 到 48 个。同时预留一个原则分区数不是越少越好因为分区数决定了未来消费者扩容的天花板。第二消费者并发数尽量等于或小于分区数。一个 12 分区的 Topic消费者并发设为 6 到 12 是比较合理的。设置超过分区数的并发是浪费因为 Kafka 不会把一个分区同时交给多个消费者消费。第三消费逻辑要设计成可重试且幂等。消息系统不能保证不重复消费消费端必须保证同一批消息处理多次和一次的效果一致。比如写数据库时用主键冲突更新或者用消息 ID 做去重表。没有幂等保护的消费逻辑在生产环境早晚会出问题。第四不要把复杂的业务处理直接写进消费者。如果消费一条消息要调用三个外部系统建议消费者里做两件事一是把消息落库或落 Redis 作为任务二是立即提交 offset然后异步任务再去处理那些耗时操作。这样可以最大限度避免消费者阻塞。比如下面的伪代码思路KafkaListener(topics order-topic, groupId order-consume-group) public void onMessage(ConsumerRecordString, String record, Acknowledgment ack) { // 快速入队后立即提交消费端保持轻量 messageTaskQueue.submit(record.value()); ack.acknowledge(); } // 异步线程池处理真实业务 Async(orderTaskExecutor) public void processTask(String message) { // 解析消息、调用服务、写库 }这种方式下Kafka 消费端永远不会积压真正的积压转移到了异步任务队列上。这时你要预留对任务队列的监控避免无限堆积。第五Kafka 积压和磁盘扩容、网盘扩容完全是两回事。很多开发者在搜索 Kafka 扩容相关话题时会看到大量和 Kafka 完全无关的磁盘扩容、网盘扩容内容。Kafka 的扩容指的是增加消费者实例、增加 Topic 分区、扩展 Broker 集群而不是给服务器加硬盘。明确这个边界能少走一些弯路。回到开头的问题遇到 Kafka 积压第一步永远不是扩机器而是回答积压发生在哪个环节、为什么消费变慢、并行度天花板在哪里。代码优化和参数调优通常比扩容更有效、更便宜、副作用更小。只有当你确认分区数充足、消费实例不足或者存在短期流量突增时扩容才是合适的应急手段。把这两类场景分清才算真正理解了 Kafka 消费模型。
返回列表