凌晨三点,Kafka 监控面板上的 Lag 数字开始不受控地往上蹿。手边第一个念头几乎都一样:消费者不够了,赶紧加实例。我在多个团队都见过这个场景,有人把消费者从 3 个扩到 10 个,结果等了十分钟,Lag 不仅没有降下来,反而因为重启和重平衡变得更难看。问题到底出在哪?
Kafka 消息堆积这件事,很少是单纯的“消费者机器不够”。它更像是一个信号,说明从生产者、Broker 到消费者再到下游系统的整条链路里,某个环节出现了失配。如果不加诊断就盲目扩容消费者,通常只是在症状上做文章,真正拖垮消费速度的地方根本没有被碰到。
1. 消息堆积,先回答“堆在哪一层”
1.1 堆积不等于消费者慢,可能只是生产太快
Kafka 的基本模型并不复杂:生产者把消息写入 Broker,Broker 按分区保存消息,消费者从分区中拉取数据并维护自己的消费位点。Lag 的本质是“分区最新位点”和“消费位点”之间的差距。这个定义意味着,Lag 涨起来不一定代表消费者出了故障,也可能是生产者瞬时写入太快,消费速度跟不上写入速度。
我在排查问题时,很少直接去看消费者。第一步是先拉一条时间线,同时看生产速率、Broker 入流量和消费速率。如果生产速率在某一时刻突然冲到每秒几万条,而消费速率仍然停留在每秒几千条,Lag 自然会上涨。这种情况下,消费者可能没有任何问题,只是它的处理能力暂时低于输入速率。
Kafka 的特点是 Broker 能扛住很高的写入,但消费者通常要做业务逻辑、访问数据库或调用外部接口,吞吐量很难和生产者直接对齐。只要写入峰值超过消费能力,堆积就会出现。这不是 Kafka 本身的故障,而是容量规划问题。
因此,看到堆积先不要默认“消费端太弱”,要问一个问题:Kafka 的 Lag 是被生产峰值推高的,还是被消费能力下降拉出来的?
1.2 三个核心指标:Lag、消费 TPS、单条耗时
Lag 本身只是一个结果指标,能反映堆积了多少,但不能回答“为什么堆积”。要定位根因,至少要同时看三个指标:
- Lag:最新位点与消费位点的差值,代表积压数量。
- 消费 TPS:消费者每秒成功消费并提交的消息条数,代表当前处理速率。
- 单条消息处理耗时:从消费者拉取到消息开始处理再到提交位点之间的时间,代表每条消息在消费链路里的延迟。
这三者的关系很直接。假设生产速率是稳定的每秒 5000 条,消费 TPS 只有 3000 条,那么即使消费者没有报错,Lag 也会稳定增长。如果消费 TPS 已经接近 5000 条,但 Lag 还在涨,就要继续看单条耗时和分区分配。
从工程经验看,真正有效的排查不是看某个时间点的 Lag 绝对值,而是看 Lag 的变化趋势与生产、消费速率是否匹配。我通常会先把这些指标放在同一张图表里对比,再决定下一步往哪个方向查。
| 指标 | 观察目的 | 堆积时常见表现 |
|---|---|---|
| Lag | 积压数量与趋势 | 持续上升,或高位不回落 |
| 消费 TPS | 消费端实时处理能力 | 显著低于生产速率 |
| 单条消息处理耗时 | 消费逻辑是否变慢 | 平均值升高,或出现长尾 |
1.3 区分“瞬时堆积”和“持续堆积”
不是所有 Lag 上涨都需要立即处理。Kafka 消费者在业务中通常允许一定的瞬时积压,尤其是面对秒杀、抢购、定时任务集中触发这类场景。瞬时堆积的特点是:生产速率在某一时间段飙升,Lag 跟着上涨,但等峰值过去后,Lag 会自然回落。
持续堆积则完全不同。它不会因为生产速率下降而自动恢复。比如消费者处理逻辑里出现了一个慢 SQL,原来单条消息 5ms 就能处理完,现在需要 500ms,消费速率急剧下降,Lag 只会越积越多。
判断这两种情况有一个简单方法:观察生产速率回落后,Lag 是否开始下降。如果下降,说明消费能力还在,只是被峰值冲垮;如果生产速率已经回到低位,Lag 依然不动甚至继续上涨,问题大概率出现在消费者或下游依赖上。
2. 为什么加消费者经常没用:四个常见误判
2.1 消费者数量超过分区数以后,加再多人也是闲置
这是很多人误会最深的一点。Kafka 的消费模型里,一个分区在同一个消费者组内,同一时刻只会分配给一个消费者。消费者可以自行选择分区去消费,但一个分区不能被组内多个消费者同时消费。
举个例子:Topic 有 10 个分区,当前消费组有 3 个消费者,每个消费者大约分配 3 到 4 个分区。如果把消费者扩到 10 个,理想情况下每个消费者处理 1 个分区,吞吐确实会提升。但如果继续扩大到 20 个,那多出来的 10 个消费者不会分配到任何分区,只能空转。
所以加消费者的第一个前置条件,是看消费者数量与分区数的关系。如果当前消费者数量已经大于或等于分区数,再加实例没有意义。我见过一些团队把消费者从 5 个加到 50 个,Topic 却只有 12 个分区,最终只有 12 个消费者在工作,剩下 38 个全部闲置,还白白增加了 Rebalance 和监控成本。
2.2 瓶颈在下游,扩容消费者只会把压力转给下游
另一种常见情况是,消费者本身 CPU、内存都不高,消费线程却大量阻塞在等待外部依赖上。比如消费消息后要写入数据库,数据库连接池被打满;或者要调用外部接口,外部接口限流导致大量重试。
这时加消费者的直接后果,是并发请求量变高,下游压力进一步增大。数据库本来就已经是瓶颈,再多几个消费者抢连接,只会让慢查询更多、锁等待更严重。消费者端的表现是处理耗时上升,Lag 不降反升。
用一句话概括:扩容消费者不会减少下游的计算量,它只会把请求更密集地打到下游。如果瓶颈在下游,先保护住下游,再谈扩容。否则就是让一个已经过载的系统继续承受更多压力。
2.3 Rebalance 风暴让消费者组一直处在重平衡中
Kafka 消费者组在成员变化时会触发 Rebalance。Rebalance 期间,所有消费者都会停止消费,等待新的分区分配方案确定。如果短时间内频繁有消费者加入、退出,或者心跳超时、处理时间超过max.poll.interval.ms,消费组就会一直处于 Rebalance 状态,看起来每一个实例都在运行,实际上却没人在消费。
盲目加消费者往往就会触发这种情况。新增实例会改变消费者组的成员关系,所有分区需要重新分配。如果实例刚启动还没完成初始化,随即发生第二次 Rebalance,或者某个消费者在处理一批消息时超过了 Max Poll Interval,消费者被移出组,又触发新的 Rebalance。这样一个循环下来,Lag 只会持续上升。
所以,不能只看消费者实例数量涨没涨,还要看消费组是否稳定。短期多次 Rebalance 是消息堆积中非常隐蔽但杀伤力极大的原因。
2.4 单分区顺序处理限制:不是简单的“人多力量大”
Kafka 只能在分区级别保证消息顺序。如果业务要求同一订单、同一用户或同一设备的消息严格按照顺序消费,那么同一个分区的消息只能由一个消费者串行处理。
单分区内的串行处理,意味着该分区吞吐上限就是这个分区内单条消息耗时的倒数。比如一条消息需要 10ms 处理,那么这个分区每秒最多处理 100 条。加多少个消费者都无法突破这个上限,因为只有当前负责该分区的消费者在消费。
要提升这种场景的吞吐,只能从业务层面改变顺序粒度。比如把单一用户的数据分散到多个分区,或者放宽顺序限制,让部分消息可以并行处理。如果业务上确实不能放宽,那么消费者数量再多也解决不了单分区瓶颈。
3. 一条能够减少误判的排查链路
3.1 先看指标:不是盯着 Lag,还看消费速率和位点提交
面对堆积,我建议按顺序排查,而不是复制粘贴网上的“加消费者”命令。
第一步,先获取消费组当前状态。Kafka 自带的命令行工具通常是第一手信息源:
kafka-consumer-groups.sh --bootstrap-server localhost:9092 \ --describe \ --group your-consumer-group这条命令会输出每个分区的current-offset、log-end-offset、lag,以及消费者实例当前分配到的分区。它至少能回答三个问题:
- 哪些分区 Lag 特别大?
- 消费者成员数和 Partition 数是否匹配?
- 是不是所有分区都有消费者持有,还是有分区处于无主状态?
只看命令输出还不够,还需要看监控里的消费 TPS 和records-lag-max。如果消费 TPS 已经接近零,说明消费者可能根本没有在拉取;如果消费 TPS 在正常水平但 Lag 不降,说明生产速率仍然高于消费速率。
3.2 再看分区分配:消费者数、分区数、消费者组是否健康
命令行输出里,CONSUMER-ID一列如果有很多空值,通常说明消费者没有分配到分区。正常的消费组,每个分区都应该有一个CURRENT-OFFSET和一个正在消费的成员。如果一个消费者组内有 20 个实例,Topic 只有 10 个分区,那么命令输出里只会看到 10 个分区,每个分区对应一个活跃消费者,另外 10 个消费者不会出现在分配列表里。
还需要留意,同一个 Topic 是否被多个消费组同时消费。不同消费组之间的消费进度互不影响,但会占 Broker 的资源。如果多个团队各自用一个新的 group 去消费同一个高频 Topic,Broker 的读压力会增加,也可能造成网络和磁盘 IO 上升。这种问题不容易直接反映在单个消费组的 Lag 上,但会表现为所有消费组同时变慢。
3.3 再看消息处理耗时:日志、链路追踪、埋点
如果在命令输出里看不到明显异常,下一步就要进到消费逻辑内部,看单条消息处理耗时。不要凭经验猜,最好在代码里加上日志或者埋点。
常见做法是在消费方法的入口和出口各记一个时间戳,输出 sample 日志:
long start = System.currentTimeMillis(); process(record); long cost = System.currentTimeMillis() - start; if (cost > 500) { log.warn("message process too slow, topic={}, partition={}, offset={}, cost={}ms", record.topic(), record.partition(), record.offset(), cost); }通过这类日志,可以判断慢的到底是哪一批消息、哪个分区、哪类业务。一般来说,单条消息处理耗时的上涨来自几个方向:
- 数据库查询变慢或连接池耗尽。
- 外部 RPC 长时间未返回。
- 消息体本身过大,导致反序列化开销高。
- JVM GC 中 Full GC 频繁。
- 消费线程遇到锁竞争。
很多时候,消息堆积的根因不是 Kafka 配置问题,而是业务逻辑里某一个依赖超时。
3.4 再看生产端:消息量、消息大小、峰值特性
如果消费者端一切正常,耗时也没有明显上涨,就要回头检查生产端。生产速率是不是出现了峰值?消息体是不是比平时大了很多?Topic 的分区数是不是长期不变,但消息量已经翻了好几倍?
生产端的问题经常是“隐性的”。比如一个上游数据库表变更,触发了大量的数据同步消息;或者一次业务迭代让单条消息从 1KB 变成了 50KB;再或者生产者重试策略有问题,把重复消息不停发给 Broker。这些问题会让 Lag 在消费者没有变化的情况下飙升。
查看生产端时,重点关注两类数据:
- 生产速率曲线:看是否有异常尖峰。
- 消息平均大小和最大大小:看是否存在大消息拖累网络和磁盘。
如果消息大小明显增加,即使每秒条数没有变化,Broker 和消费者之间的数据量也会变大,消费耗时自然上升。
4. 针对瓶颈的解法:不换方案,先解瓶颈
4.1 消费逻辑本身慢:批处理、异步化、并行化
如果确认瓶颈在消费业务逻辑本身,比如每条消息都要执行数据库写入、查一次用户信息、调用一次远程接口,那么首先要优化的不是消费者数量,而是消费方式。
一种常见手段是批量处理。Kafka 消费者可以通过poll一次拉取多条消息,在业务允许的情况下,把它们聚合成一个批次,再一次性写数据库或批量调用外部接口。这样能把多次网络交互合并成一次,吞吐提升会很显著。
另一种手段是异步化。消费者线程只负责拉取消息和投递到线程池,真正耗时的业务逻辑由线程池并发执行。不过异步化要非常注意提交位点的时机,如果消息已经投递到线程池但还没有处理完,就提交了位点,一旦进程崩溃就会出现消息丢失。推荐的做法是等一批消息处理完成后,再调用commitSync或commitAsync。
这里要提醒一下:批处理和异步化并不总是适用。如果业务严格要求顺序,异步化很可能破坏顺序。这种情况通常要考虑先优化单条消息处理链路,减少数据库查询次数和外部调用次数,降低单条耗时。
4.2 分区分配不均:key 散列、分区数规划、消费者数量匹配
有时候消费者总数和分区数都没有问题,但消息的 key 分布严重不均,导致某些分区堆积明显,另一些分区却很空闲。
比如订单消息用店铺 ID 做 key,而某个大店铺的消息量占了一半,这个店铺所在分区的消费者就会成为热点。就算整体消费速率看起来正常,热点分区依然会有 Lag。
这种情况下,加消费者没有用,因为热点只有一个分区,加再多消费者也无法分担该分区的消息。需要做的是让热点 key 能够分散到更多分区。具体方法包括:
- 更换 key 设计,使用更均匀的业务标识。
- 对热点 key 加随机后缀,但前提是业务允许乱序。
- 调整分区策略,在生产者端自定义分区器。
- 如果顺序要求严格,可以考虑拆分 Topic,把热点业务拆分到独立 Topic,并配置更多分区。
分区数本身也需要提前规划。Topic 一旦创建,分区数在早期还可以调整,但分区数扩容会导致同一个 key 的消息可能被分到不同分区,影响顺序。所以创建高吞吐 Topic 时,最好先根据峰值消息量估算出合理分区数,而不是用一个小分区数去扛大流量。
4.3 下游瓶颈:限流、降级、异步削峰、扩容外部依赖
如果排查到最后发现瓶颈是数据库连接池、外部 API 限流或者下游系统处理不过来,那么消费者实例数量越加,下游会越危险。
正确的思路是先给消费端加保护。可以在消费者处理逻辑里加入信号量或限流器,控制对下游的并发请求数量。虽然这会暂时让 Lag 保持高位,但至少不会拖垮下游系统。下游系统一旦恢复稳定,再逐步放开消费速度。
另一种思路是把“消费”拆成两个阶段。第一阶段,消费者从 Kafka 拉取消息后,立刻写入一个内部缓冲队列或者本地表,不直接执行业务逻辑;第二阶段,由独立的 worker 线程池从这个缓冲里取数据,按下游能承受的速度处理。这样 Kafka 端只负责快速拉取和提交,业务处理被削峰填谷。
如果外部依赖确实需要扩容量,那也要在影响评估之后再执行。比如数据库连接池从 50 调到 200,要考虑数据库本身的 CPU、内存和连接数上限,不能简单加大连接池。扩容下游永远比扩容消费者更接近根因。
4.4 什么时候加消费者才有效:三个前置条件
说了这么多,并不是说“加消费者”完全不可行。它只是不是第一优先级。如果一个场景同时满足以下几个条件,扩容消费者通常是有效且成本最低的方案:
- 当前消费者数量小于 Topic 分区数,且分区数还有比较大的余量。
- 消费者实例自身的 CPU、内存、网络、磁盘都没有达到瓶颈。
- 已确认消费端处理逻辑和下游依赖不是主要瓶颈,瓶颈只是“固定并发数太低”。
比如用 Spring Boot 写了一个消费逻辑,逻辑很简单,只是把消息转发到本地日志文件,没有外部依赖。这台消费者单实例处理能力有限,而 Topic 有 32 个分区,当前只有 4 个消费者,每个消费者分配 8 个分区。这种情况加消费者是完全合理的,每加一个实例,整体吞吐都会上升。
但在加之前,还是要确认新增实例不会引发频繁 Rebalance。尽量在低峰期操作,或者使用静态成员分组,减少 Rebalance 带来的停顿。
5. 把“救火”变成“预防”:给团队的落地框架
5.1 建立基线:先知道正常水位
很多团队处理堆积时很被动,是因为不知道“正常”是什么样子。平时没有记录消费组消费速率和 Lag 的基线,等到报警时只能靠猜。
我建议每个消费组都建立一份简单的水位记录,至少包括:正常情况下消费者的数量、每个消费者的处理能力、生产速率峰值、Lag 的安全阈值。有了这些数据,才能判断当前 Lag 是异常还是季节性波动。
基线的建立不复杂。从监控系统里拉出最近一个月的数据,取 P95 和 P99 值作为参考。再结合业务高峰期,比如每天几点集中推送,就能形成一张消费组健康度表。
5.2 设置分级告警和自动诊断
不要等到 Lag 涨到几百万才人工介入。告警应该分级别:
- 提示级:Lag 超过安全阈值的 1.5 倍,持续 5 分钟,记录事件。
- 警告级:Lag 持续上升且消费速率低于生产速率,发送通知。
- 严重级:Lag 导致业务结果延迟超过可接受范围,触发应急预案。
告警消息里不要只写“消费组 lag 超限”。最好附带自动诊断信息:当前分区数、消费者数、消费 TPS、生产 TPS、最近是否有 Rebalance。这样一来,值班的人不用登录服务器查半天,就能判断大概方向。
要做到这一步,通常需要一个脚本或定时任务,调用kafka-consumer-groups.sh获取状态,再通过监控 API 取生产和消费速率,把信息拼成一条告警发出去。初期可以做得简单,但一定要有这个自动化过程。
5.3 压测与容量规划
Kafka 消费端的容量规划,不像加内存加 CPU 那么简单。建议在预发环境做一次消费压测,模拟高峰期消息量,测量单个消费者在单个分区上的真实处理速率。
有一个简单估算公式可以参考:
需要的消费者数 ≈ 生产峰值消息量 / (单分区消费速率 × 分区数)这里的分区数指的是该 Topic 分配给这个消费组的分区数。如果算出来的消费者数超过分区数,说明要么分配不合理,要么单分区处理太慢。压测的目的不是得到精确公式,而是弄清楚“单条消息处理耗时”和“单消费者处理能力”这两个关键数字。
压测之后,还要给消费系统留出 30% 到 50% 的容量冗余。线上不可能永远按预发环境的参数运行,网络抖动、GC、下游慢查询都会导致实际吞吐下降。
5.4 从一次堆积复盘到长期改进
每次堆积被恢复后,最该做的是复盘,而不是庆祝。建议记录一份简单的事件复盘,包含:
- 堆积发生的时间、持续时长、最大 Lag。
- 当时的消费组指标和生产端指标截图。
- 根因是什么:生产峰值、下游瓶颈、消费者逻辑、还是 Rebalance。
- 采用了什么恢复动作:限流、扩容、修改参数、优化代码。
- 后续需要做的改进项。
长期来看,Kafka 消息堆积不是一个可以用加机器一劳永逸解决的问题。它考验的是团队对系统容量的理解、对消费逻辑的掌控、对监控告警的建设,还有面对报警时能不能先冷静判断,再动手操作。
下次看到 Lag 上涨,可以把“加消费者”这个念头往后放一放,先打开监控,拉出生产速率和消费速率,数一数分区数和消费者数,看一眼消费耗时的变化曲线。大多数堆积的答案,都在这些数据里,不在扩容按钮上。