MQ Lag 分布与堆积定位
面试里问“Lag 分布是什么”,不是问你会不会算总 Lag,而是看你是否真的懂 MQ 堆积定位。总 Lag 只能说明“积压很多”,不能说明“为什么积压”。真正排查时要看每个分区、队列、MessageQueue、消费者通道分别积压多少,这就是 Lag 分布。
这一页解决这些问题
| 问题 | 要掌握到什么程度 |
|---|---|
| Lag 是什么 | 知道生产进度和消费进度之间的差值 |
| Lag 分布是什么 | 知道每个分区或队列分别落后多少 |
| 为什么不能只看总 Lag | 能区分整体消费慢和局部热点卡住 |
| Kafka 怎么看 | 知道 Partition Lag、Consumer Group、分区并行度 |
| RocketMQ 怎么看 | 知道 MessageQueue Diff、消费组进度 |
| RabbitMQ 怎么类比 | 知道 Ready、Unacked、Queue 维度分布 |
| 怎么处理 | 能按均匀堆积、倾斜堆积、单通道卡死、Lag 波动分别定位 |
Lag 是什么
Lag 可以理解为:消费者已经处理到的位置,落后于 Broker 最新消息位置多少。
以 Kafka 为例:
Lag = LOG-END-OFFSET - CURRENT-OFFSET| 字段 | 含义 |
|---|---|
CURRENT-OFFSET | 当前消费组已经提交到的位置 |
LOG-END-OFFSET | 当前分区最新消息位置 |
LAG | 当前分区还有多少消息没被这个消费组处理 |
示例:
TOPIC PARTITION CURRENT-OFFSET LOG-END-OFFSET LAG
order-event 0 1000 1500 500表示 order-event 的 partition 0 里,这个消费组还落后 500 条消息。
Lag 分布是什么
Lag 分布是看每个分区或队列分别落后多少,而不是只看总和。
flowchart TD
A["Consumer Group 总 Lag"] --> B["Partition 0 Lag"]
A --> C["Partition 1 Lag"]
A --> D["Partition 2 Lag"]
A --> E["Partition 3 Lag"]
B --> F["判断是否均匀"]
C --> F
D --> F
E --> F例如总 Lag 都是 12000,但下面两种情况完全不同。
均匀分布说明什么
Partition 0 Lag = 3000
Partition 1 Lag = 3000
Partition 2 Lag = 3000
Partition 3 Lag = 3000
Total Lag = 12000这说明多个分区都在落后,更像是整体消费能力不足,或者所有消费者共同依赖的下游系统变慢。
常见原因:
- 消费者实例或线程整体不足。
- 数据库、ES、Redis、HTTP 接口整体变慢。
- Broker 拉取或投递能力下降。
- 生产 TPS 突然超过容量预案。
- 消费逻辑版本发布后整体变慢。
处理方向:
| 方向 | 说明 |
|---|---|
| 扩容消费者 | 只在分区/队列并行度允许且下游扛得住时有效 |
| 提高单消费者能力 | 批量消费、减少单条耗时、优化反序列化 |
| 优化下游 | 慢 SQL、ES Bulk、Redis、HTTP 连接池 |
| 限流生产 | 防止堆积继续扩大 |
| 消化存量 | 计算净消化速度,必要时临时扩容 |
倾斜分布说明什么
Partition 0 Lag = 100
Partition 1 Lag = 11600
Partition 2 Lag = 150
Partition 3 Lag = 150
Total Lag = 12000这说明堆积几乎集中在 Partition 1,更像是局部问题。
常见原因:
- 热点 key 都被路由到同一个分区。
- 某条慢消息卡住了分区内顺序处理。
- 分配到这个分区的消费者实例异常。
- 该分区消息对应的业务数据特别慢,比如同一个大客户、大订单、大文件。
- 该分区一直重试失败,正常消息被拖住。
处理方向:
| 方向 | 说明 |
|---|---|
| 定位最大 Lag 分区 | 找到具体 Partition、MessageQueue 或 Queue |
| 定位负责消费者 | 看该消费者日志、线程栈、CPU、GC、线程池 |
| 抽样消息 key | 判断是否热点 key 或大客户流量 |
| 查慢消息 | 看是否单条消息处理很慢或反复失败 |
| 失败隔离 | 对不可重试错误进入死信,避免阻塞正常消息 |
为什么不能只看总 Lag
flowchart TD
A["发现总 Lag 很高"] --> B{"看 Lag 分布"}
B -- "各分区都高" --> C["整体能力不足"]
C --> D["扩容消费者、优化下游、限流生产"]
B -- "少数分区高" --> E["局部热点或阻塞"]
E --> F["查热点 key、慢消息、异常实例、顺序阻塞"]| 只看总 Lag 的误判 | 真相可能是 |
|---|---|
| 以为消费者都慢 | 实际只有一个分区慢 |
| 以为加消费者能解决 | 实际分区数不足或热点 key 卡住 |
| 以为 Broker 有问题 | 实际是某个消费者线程死锁 |
| 以为可以清空队列 | 实际是少量异常消息导致局部阻塞 |
| 以为 Lag 不高就安全 | 实际最老消息已经超 SLA |
总 Lag 看规模,Lag 分布看病灶。
Kafka 的 Lag 分布怎么看
命令:
kafka-consumer-groups.sh \
--bootstrap-server 127.0.0.1:9092 \
--describe \
--group order-search-sync-group示例输出:
TOPIC PARTITION CURRENT-OFFSET LOG-END-OFFSET LAG CONSUMER-ID
order-event 0 10000 10100 100 consumer-a
order-event 1 5000 16000 11000 consumer-b
order-event 2 9900 10050 150 consumer-c阅读顺序:
- 先看总 Lag 是否超过业务阈值。
- 再看每个 Partition 的 Lag 是否均匀。
- 找出 Lag 最大的 Partition。
- 看它分配给哪个 Consumer。
- 查该 Consumer 的日志、线程池、处理耗时、错误率。
- 查该 Partition 的 key 分布和慢消息。
flowchart TD
A["Kafka Consumer Lag"] --> B["按 Partition 排序"]
B --> C["找到 Lag 最大分区"]
C --> D["定位负责该分区的 Consumer"]
D --> E["查消费者日志和线程栈"]
C --> F["抽样该分区消息 key"]
F --> G["判断热点 key 或异常消息"]Kafka 同一个 Consumer Group 中,一个 Partition 同一时刻只能被一个消费者消费。
flowchart TD
A["Topic 4 个 Partition"] --> B["最多 4 个消费者并行消费"]
B --> C{"启动 10 个消费者"}
C -- "4 个有分区" --> D["真正工作"]
C -- "6 个无分区" --> E["空闲,不能提升吞吐"]所以消费者实例数已经等于或超过分区数时,继续加实例通常不能提升吞吐。需要看是否增加分区、提升单消费者能力、改分区 key 或优化下游。
RocketMQ 的 Lag 分布怎么看
RocketMQ 里常用 ConsumerGroup 的消费进度和 Diff 判断堆积。它的思想和 Kafka 类似:不要只看总 Diff,要看每个 MessageQueue 的 Diff。
MessageQueue Diff = Broker 最大位点 - ConsumerGroup 消费位点
Total Diff = 所有 MessageQueue Diff 之和排查时关注:
| 维度 | 看什么 |
|---|---|
| Topic | 哪类业务消息堆积 |
| ConsumerGroup | 哪个消费组落后 |
| MessageQueue | 是否只有某些队列 Diff 高 |
| Retry Topic | 是否大量失败消息重试 |
| DLQ | 是否死信增长 |
| Consumer 实例 | 是否某个实例消费慢、频繁重启 |
flowchart TD
A["RocketMQ Diff 高"] --> B{"是否所有 MessageQueue 都高"}
B -- "是" --> C["整体消费能力不足或下游慢"]
B -- "否" --> D["局部队列热点或顺序阻塞"]
C --> E["看消费线程、DB/ES/HTTP、生产 TPS"]
D --> F["查 MessageQueue、消息 Key、慢消息、重试"]RocketMQ 顺序消息尤其要注意:同一个队列内需要顺序处理,如果某条消息一直失败,后面的消息可能被阻塞。此时盲目扩容消费者不会让这个队列并行起来。
RabbitMQ 怎么类比 Lag 分布
RabbitMQ 没有 Kafka 那种 offset 模型,但也要看“分布”,不要只看总消息数。
| 指标 | 含义 |
|---|---|
| Ready | 队列中等待投递的消息 |
| Unacked | 已投递给消费者但未确认的消息 |
| Consumers | 消费者数量 |
| Publish rate | 生产速率 |
| Deliver rate | 投递速率 |
| Ack rate | 确认速率 |
RabbitMQ 的“分布”主要看:
- 哪些 Queue 的 Ready 高。
- 哪些 Queue 的 Unacked 高。
- 同一业务是否拆了多个 Queue,其中某个 Queue 特别高。
- 某些消费者是否持有大量 Unacked。
- prefetch 是否导致消息被少数消费者拿走。
flowchart TD
A["RabbitMQ 堆积"] --> B{"Ready 还是 Unacked 高"}
B -- "Ready 高" --> C["消息在队列里等投递"]
C --> D["看消费者数量、prefetch、投递速率"]
B -- "Unacked 高" --> E["消息已被消费者拿走"]
E --> F["看业务耗时、ACK、消费者线程"]典型判断:
| 现象 | 说明 |
|---|---|
| Ready 高,Unacked 低 | 消费者不够、没拉取、prefetch 太小、消费者掉线 |
| Ready 低,Unacked 高 | 消费者拿了但处理慢或不 ACK |
| Ready 和 Unacked 都高 | 生产过快,消费整体处理不过来 |
| 单个 Queue 高,其他正常 | 路由热点、队列消费者异常、该业务下游慢 |
Lag 分布和热点 key
热点 key 是 Lag 倾斜最常见原因之一。
假设 Kafka 按 userId 分区:
partition = hash(userId) % partitionCount如果某个大客户 userId=1001 产生了大量消息,它们会集中到同一个分区。
flowchart TD
A["大量 userId=1001 消息"] --> B["hash 到同一 Partition"]
B --> C["该 Partition 消费压力远高于其他分区"]
C --> D["Lag 倾斜"]处理方向:
| 方案 | 适用场景 | 风险 |
|---|---|---|
| 调整分区 key | key 设计不合理 | 可能影响同 key 顺序 |
| 对热点 key 二次拆分 | 热点业务可拆并行 | 需要消费端聚合或保证局部顺序 |
| 热点单独 Topic | 大客户、大租户、高频业务 | 架构复杂度上升 |
| 批处理热点消息 | 下游支持批量写 | 延迟和失败处理更复杂 |
| 限流热点来源 | 上游可控 | 需要业务接受降速 |
如果业务强依赖同一个订单、同一个用户的顺序,就不能简单把 key 打散。要在顺序和吞吐之间做权衡。
Lag 分布和慢消息
有时不是热点 key,而是一条或一批慢消息卡住消费。
| 慢消息类型 | 示例 |
|---|---|
| 大消息 | 消息体很大,反序列化或网络传输慢 |
| 脏数据 | 字段缺失、格式错误,反复异常重试 |
| 下游慢数据 | 某个订单关联数据特别多,SQL 慢 |
| 外部接口慢 | 某条消息触发第三方接口超时 |
| 锁冲突消息 | 多条消息竞争同一业务资源 |
排查步骤:
flowchart TD
A["某分区 Lag 不动"] --> B["定位负责消费者"]
B --> C["查当前正在处理的消息 key"]
C --> D["查消费日志耗时"]
D --> E{"是否单条消息反复失败"}
E -- "是" --> F["失败隔离、死信、人工修复"]
E -- "否" --> G["查下游慢 SQL、接口、锁等待"]Lag 分布和重试风暴
重试风暴是指失败消息反复被快速重试,占满消费者处理能力,导致正常消息也处理不过来。
flowchart TD
A["某批消息消费失败"] --> B["进入重试"]
B --> C["很快再次投递"]
C --> D["再次失败"]
D --> E["消费者大量时间处理失败消息"]
E --> F["正常消息消费 TPS 下降"]
F --> G["Lag 继续增长"]判断信号:
- 错误率突然升高。
- Retry Topic 或重试队列增长。
- DLQ 开始增长。
- 消费 TPS 看起来不低,但成功 TPS 很低。
- 日志里同一类异常重复出现。
处理要点:
| 动作 | 说明 |
|---|---|
| 区分可重试和不可重试 | 参数错误、格式错误不要无限重试 |
| 延迟重试 | 避免立即重试打爆消费者 |
| 死信隔离 | 超过次数进入 DLQ,保护正常消息 |
| 告警和人工处理 | 死信不是垃圾桶,必须有人处理 |
| 幂等 | 重试一定会带来重复消费风险 |
扩容后再次堆积,怎么结合 Lag 分布判断
flowchart TD
A["扩容后再次堆积"] --> B["看 Lag 分布变化"]
B --> C{"扩容后是否所有分区都下降"}
C -- "是,后面又全部上涨" --> D["整体下游瓶颈被打出来"]
C -- "否,只有少数分区不降" --> E["热点分区或慢消息"]
C -- "总 Lag 降但最老延迟高" --> F["有长尾消息或顺序阻塞"]
D --> G["查 DB/ES/HTTP、连接池、线程池"]
E --> H["查 key 分布、负责消费者、重试"]
F --> I["查最老消息和单条耗时"]| 分布变化 | 说明 |
|---|---|
| 所有分区 Lag 先降后升 | 扩容把整体下游瓶颈打出来了 |
| 大多数分区下降,个别不降 | 个别分区热点或卡死 |
| 总 Lag 下降,但业务仍超时 | 最老消息延迟高,可能顺序阻塞 |
| Lag 频繁归零又暴涨 | 消费者重启、Rebalance、生产批量尖峰 |
商业场景:订单同步 ES 堆积
订单服务写 MySQL 后,通过 MQ 同步 ES 搜索索引。
flowchart TD
A["订单状态变化"] --> B["写 MySQL"]
B --> C["发送订单事件到 MQ"]
C --> D["搜索同步消费者"]
D --> E["查 MySQL 最新快照"]
E --> F["Bulk 写 ES"]
F --> G["提交 Offset 或 ACK"]如果 Lag 均匀升高:
- 可能 ES 整体写入慢。
- 可能 Bulk 太小或太大。
- 可能消费者线程池不足。
- 可能 MySQL 查询最新快照慢。
如果只有少数分区 Lag 高:
- 可能某些大客户订单量集中。
- 可能某个订单状态消息反复失败。
- 可能某个消费者实例连接 ES 异常。
- 可能分区 key 设计导致热点。
处理时不能直接 reset offset 到 latest,因为这会导致 ES 缺数据。正确做法是先限流或暂停补数据,再定位慢点,失败消息进入死信或异常表,最后用 MySQL 主库补偿重建 ES。
可运行 Demo:计算 Lag 分布摘要
下面这个 Demo 不依赖具体 MQ SDK,用一组分区 Lag 数据演示如何判断“均匀堆积”还是“倾斜堆积”。真实项目里数据来自 Kafka、RocketMQ 或监控系统。
import java.util.Comparator;
import java.util.List;
public class LagDistributionDemo {
record PartitionLag(String name, long lag) {}
public static void main(String[] args) {
List<PartitionLag> lags = List.of(
new PartitionLag("partition-0", 120),
new PartitionLag("partition-1", 11600),
new PartitionLag("partition-2", 180),
new PartitionLag("partition-3", 100)
);
long total = lags.stream().mapToLong(PartitionLag::lag).sum();
PartitionLag max = lags.stream()
.max(Comparator.comparingLong(PartitionLag::lag))
.orElseThrow();
double maxRatio = total == 0 ? 0 : (max.lag() * 1.0 / total);
System.out.println("totalLag=" + total);
System.out.println("maxLagPartition=" + max.name() + ", lag=" + max.lag());
System.out.printf("maxRatio=%.2f%%%n", maxRatio * 100);
if (maxRatio > 0.7) {
System.out.println("判断:Lag 高度集中,优先排查热点 key、慢消息、单消费者异常。");
} else {
System.out.println("判断:Lag 分布较均匀,优先排查整体消费能力和下游瓶颈。");
}
}
}输出类似:
totalLag=12000
maxLagPartition=partition-1, lag=11600
maxRatio=96.67%
判断:Lag 高度集中,优先排查热点 key、慢消息、单消费者异常。这个思路可以直接落到监控告警:
| 指标 | 含义 |
|---|---|
total_lag | 总积压规模 |
max_partition_lag | 最大单分区积压 |
max_lag_ratio | 最大分区 Lag / 总 Lag |
oldest_message_age | 最老消息等待时间 |
success_tps | 真正成功处理的 TPS |
retry_rate | 重试比例 |
排查清单
| 步骤 | 要回答的问题 |
|---|---|
| 1 | 总 Lag 是多少,是否超过 SLA |
| 2 | 最老消息等待多久,比总 Lag 更能反映业务风险 |
| 3 | Lag 是均匀分布还是集中在少数分区/队列 |
| 4 | 最大 Lag 的分区由哪个消费者处理 |
| 5 | 该消费者是否 CPU、线程池、连接池、GC、日志异常 |
| 6 | 最大 Lag 分区里的 key 是否集中 |
| 7 | 是否有慢消息、脏消息、重复失败消息 |
| 8 | Retry、DLQ、错误率是否同步上升 |
| 9 | 下游 DB、ES、Redis、HTTP 是否变慢 |
| 10 | 扩容是否受分区/队列数限制 |
常见坑
| 坑 | 后果 | 正确做法 |
|---|---|---|
| 只看总 Lag | 不知道整体慢还是局部卡 | 看分区、队列、消费者维度分布 |
| Lag 高就加消费者 | 分区不足、热点、下游慢时无效 | 先判断瓶颈 |
| Lag 降了就认为恢复 | 最老消息可能仍超 SLA | 同时看最老消息等待时间 |
| 直接 reset offset | 未处理消息被跳过,业务数据缺失 | 只有可丢弃或可重建时才考虑 |
| 失败消息无限重试 | 正常消息被拖慢,重试风暴 | 延迟重试、死信、人工修复 |
| 忽略 key 分布 | 热点长期存在 | 设计合理分区 key 和热点拆分策略 |
面试标准回答
Lag 表示消费进度落后生产进度多少。以 Kafka 为例,单分区 Lag 等于 LOG-END-OFFSET 减 CURRENT-OFFSET。Lag 分布不是只看总 Lag,而是看每个分区、队列或 MessageQueue 分别落后多少。
只看总 Lag 只能知道积压规模,不能定位原因。如果所有分区 Lag 都高,通常说明整体消费能力不足、下游 DB/ES/接口慢、生产 TPS 超过容量;如果只有少数分区 Lag 高,通常说明热点 key、慢消息、单消费者异常、顺序消息阻塞或重试风暴。
排查时我会先看生产 TPS、消费 TPS、总 Lag、最老消息等待时间,再看 Lag 分布。Kafka 看每个 Partition Lag 和负责的 Consumer;RocketMQ 看每个 MessageQueue 的 Diff;RabbitMQ 虽然没有 offset Lag,但要看各 Queue 的 Ready、Unacked、Consumers、Ack rate。处理上不能只靠加消费者,要结合分区并行度、下游瓶颈、热点 key、失败重试和业务是否允许补偿来决定。关联知识点
| 知识点 | 说明 |
|---|---|
| 消息堆积与背压 | 堆积总览、背压、排查流程 |
| 消费者扩容后再次堆积 | 扩容短暂有效后再次堆积的原因 |
| Kafka 分区与 Offset | Kafka 分区、顺序和 Offset 原理 |
| RocketMQ 顺序消息 | 队列顺序和阻塞问题 |
| RabbitMQ ACK 重试和死信 | Ready、Unacked、ACK 和死信 |
