Skip to content

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 为例:

text
Lag = LOG-END-OFFSET - CURRENT-OFFSET
字段含义
CURRENT-OFFSET当前消费组已经提交到的位置
LOG-END-OFFSET当前分区最新消息位置
LAG当前分区还有多少消息没被这个消费组处理

示例:

text
TOPIC        PARTITION  CURRENT-OFFSET  LOG-END-OFFSET  LAG
order-event  0          1000            1500            500

表示 order-eventpartition 0 里,这个消费组还落后 500 条消息。

Lag 分布是什么

Lag 分布是看每个分区或队列分别落后多少,而不是只看总和。

mermaid
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,但下面两种情况完全不同。

均匀分布说明什么

text
Partition 0 Lag = 3000
Partition 1 Lag = 3000
Partition 2 Lag = 3000
Partition 3 Lag = 3000
Total Lag       = 12000

这说明多个分区都在落后,更像是整体消费能力不足,或者所有消费者共同依赖的下游系统变慢。

常见原因:

  1. 消费者实例或线程整体不足。
  2. 数据库、ES、Redis、HTTP 接口整体变慢。
  3. Broker 拉取或投递能力下降。
  4. 生产 TPS 突然超过容量预案。
  5. 消费逻辑版本发布后整体变慢。

处理方向:

方向说明
扩容消费者只在分区/队列并行度允许且下游扛得住时有效
提高单消费者能力批量消费、减少单条耗时、优化反序列化
优化下游慢 SQL、ES Bulk、Redis、HTTP 连接池
限流生产防止堆积继续扩大
消化存量计算净消化速度,必要时临时扩容

倾斜分布说明什么

text
Partition 0 Lag = 100
Partition 1 Lag = 11600
Partition 2 Lag = 150
Partition 3 Lag = 150
Total Lag       = 12000

这说明堆积几乎集中在 Partition 1,更像是局部问题。

常见原因:

  1. 热点 key 都被路由到同一个分区。
  2. 某条慢消息卡住了分区内顺序处理。
  3. 分配到这个分区的消费者实例异常。
  4. 该分区消息对应的业务数据特别慢,比如同一个大客户、大订单、大文件。
  5. 该分区一直重试失败,正常消息被拖住。

处理方向:

方向说明
定位最大 Lag 分区找到具体 Partition、MessageQueue 或 Queue
定位负责消费者看该消费者日志、线程栈、CPU、GC、线程池
抽样消息 key判断是否热点 key 或大客户流量
查慢消息看是否单条消息处理很慢或反复失败
失败隔离对不可重试错误进入死信,避免阻塞正常消息

为什么不能只看总 Lag

mermaid
flowchart TD
    A["发现总 Lag 很高"] --> B{"看 Lag 分布"}
    B -- "各分区都高" --> C["整体能力不足"]
    C --> D["扩容消费者、优化下游、限流生产"]
    B -- "少数分区高" --> E["局部热点或阻塞"]
    E --> F["查热点 key、慢消息、异常实例、顺序阻塞"]
只看总 Lag 的误判真相可能是
以为消费者都慢实际只有一个分区慢
以为加消费者能解决实际分区数不足或热点 key 卡住
以为 Broker 有问题实际是某个消费者线程死锁
以为可以清空队列实际是少量异常消息导致局部阻塞
以为 Lag 不高就安全实际最老消息已经超 SLA

总 Lag 看规模,Lag 分布看病灶。

Kafka 的 Lag 分布怎么看

命令:

bash
kafka-consumer-groups.sh \
  --bootstrap-server 127.0.0.1:9092 \
  --describe \
  --group order-search-sync-group

示例输出:

text
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

阅读顺序:

  1. 先看总 Lag 是否超过业务阈值。
  2. 再看每个 Partition 的 Lag 是否均匀。
  3. 找出 Lag 最大的 Partition。
  4. 看它分配给哪个 Consumer。
  5. 查该 Consumer 的日志、线程池、处理耗时、错误率。
  6. 查该 Partition 的 key 分布和慢消息。
mermaid
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 同一时刻只能被一个消费者消费。

mermaid
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。

text
MessageQueue Diff = Broker 最大位点 - ConsumerGroup 消费位点
Total Diff = 所有 MessageQueue Diff 之和

排查时关注:

维度看什么
Topic哪类业务消息堆积
ConsumerGroup哪个消费组落后
MessageQueue是否只有某些队列 Diff 高
Retry Topic是否大量失败消息重试
DLQ是否死信增长
Consumer 实例是否某个实例消费慢、频繁重启
mermaid
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 的“分布”主要看:

  1. 哪些 Queue 的 Ready 高。
  2. 哪些 Queue 的 Unacked 高。
  3. 同一业务是否拆了多个 Queue,其中某个 Queue 特别高。
  4. 某些消费者是否持有大量 Unacked。
  5. prefetch 是否导致消息被少数消费者拿走。
mermaid
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 分区:

text
partition = hash(userId) % partitionCount

如果某个大客户 userId=1001 产生了大量消息,它们会集中到同一个分区。

mermaid
flowchart TD
    A["大量 userId=1001 消息"] --> B["hash 到同一 Partition"]
    B --> C["该 Partition 消费压力远高于其他分区"]
    C --> D["Lag 倾斜"]

处理方向:

方案适用场景风险
调整分区 keykey 设计不合理可能影响同 key 顺序
对热点 key 二次拆分热点业务可拆并行需要消费端聚合或保证局部顺序
热点单独 Topic大客户、大租户、高频业务架构复杂度上升
批处理热点消息下游支持批量写延迟和失败处理更复杂
限流热点来源上游可控需要业务接受降速

如果业务强依赖同一个订单、同一个用户的顺序,就不能简单把 key 打散。要在顺序和吞吐之间做权衡。

Lag 分布和慢消息

有时不是热点 key,而是一条或一批慢消息卡住消费。

慢消息类型示例
大消息消息体很大,反序列化或网络传输慢
脏数据字段缺失、格式错误,反复异常重试
下游慢数据某个订单关联数据特别多,SQL 慢
外部接口慢某条消息触发第三方接口超时
锁冲突消息多条消息竞争同一业务资源

排查步骤:

mermaid
flowchart TD
    A["某分区 Lag 不动"] --> B["定位负责消费者"]
    B --> C["查当前正在处理的消息 key"]
    C --> D["查消费日志耗时"]
    D --> E{"是否单条消息反复失败"}
    E -- "是" --> F["失败隔离、死信、人工修复"]
    E -- "否" --> G["查下游慢 SQL、接口、锁等待"]

Lag 分布和重试风暴

重试风暴是指失败消息反复被快速重试,占满消费者处理能力,导致正常消息也处理不过来。

mermaid
flowchart TD
    A["某批消息消费失败"] --> B["进入重试"]
    B --> C["很快再次投递"]
    C --> D["再次失败"]
    D --> E["消费者大量时间处理失败消息"]
    E --> F["正常消息消费 TPS 下降"]
    F --> G["Lag 继续增长"]

判断信号:

  1. 错误率突然升高。
  2. Retry Topic 或重试队列增长。
  3. DLQ 开始增长。
  4. 消费 TPS 看起来不低,但成功 TPS 很低。
  5. 日志里同一类异常重复出现。

处理要点:

动作说明
区分可重试和不可重试参数错误、格式错误不要无限重试
延迟重试避免立即重试打爆消费者
死信隔离超过次数进入 DLQ,保护正常消息
告警和人工处理死信不是垃圾桶,必须有人处理
幂等重试一定会带来重复消费风险

扩容后再次堆积,怎么结合 Lag 分布判断

mermaid
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 搜索索引。

mermaid
flowchart TD
    A["订单状态变化"] --> B["写 MySQL"]
    B --> C["发送订单事件到 MQ"]
    C --> D["搜索同步消费者"]
    D --> E["查 MySQL 最新快照"]
    E --> F["Bulk 写 ES"]
    F --> G["提交 Offset 或 ACK"]

如果 Lag 均匀升高:

  1. 可能 ES 整体写入慢。
  2. 可能 Bulk 太小或太大。
  3. 可能消费者线程池不足。
  4. 可能 MySQL 查询最新快照慢。

如果只有少数分区 Lag 高:

  1. 可能某些大客户订单量集中。
  2. 可能某个订单状态消息反复失败。
  3. 可能某个消费者实例连接 ES 异常。
  4. 可能分区 key 设计导致热点。

处理时不能直接 reset offset 到 latest,因为这会导致 ES 缺数据。正确做法是先限流或暂停补数据,再定位慢点,失败消息进入死信或异常表,最后用 MySQL 主库补偿重建 ES。

可运行 Demo:计算 Lag 分布摘要

下面这个 Demo 不依赖具体 MQ SDK,用一组分区 Lag 数据演示如何判断“均匀堆积”还是“倾斜堆积”。真实项目里数据来自 Kafka、RocketMQ 或监控系统。

java
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 分布较均匀,优先排查整体消费能力和下游瓶颈。");
        }
    }
}

输出类似:

text
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 更能反映业务风险
3Lag 是均匀分布还是集中在少数分区/队列
4最大 Lag 的分区由哪个消费者处理
5该消费者是否 CPU、线程池、连接池、GC、日志异常
6最大 Lag 分区里的 key 是否集中
7是否有慢消息、脏消息、重复失败消息
8Retry、DLQ、错误率是否同步上升
9下游 DB、ES、Redis、HTTP 是否变慢
10扩容是否受分区/队列数限制

常见坑

后果正确做法
只看总 Lag不知道整体慢还是局部卡看分区、队列、消费者维度分布
Lag 高就加消费者分区不足、热点、下游慢时无效先判断瓶颈
Lag 降了就认为恢复最老消息可能仍超 SLA同时看最老消息等待时间
直接 reset offset未处理消息被跳过,业务数据缺失只有可丢弃或可重建时才考虑
失败消息无限重试正常消息被拖慢,重试风暴延迟重试、死信、人工修复
忽略 key 分布热点长期存在设计合理分区 key 和热点拆分策略

面试标准回答

text
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 分区与 OffsetKafka 分区、顺序和 Offset 原理
RocketMQ 顺序消息队列顺序和阻塞问题
RabbitMQ ACK 重试和死信Ready、Unacked、ACK 和死信