概述

凌晨三点,手机被告警轰炸。Kafka 某个消费者组的 Lag 从 200 飙到 50 万,下游实时报表全部卡住,数据团队在群里疯狂 at 你。你爬起来连上跳板机,敲 kafka-consumer-groups.sh --describe,看到 Lag 列一串刺眼的数字。接下来半小时,你要回答三个问题:积压在哪个分区?消费者是死了还是慢了?加机器能救吗?

消息队列是分布式系统的"下水道"——平时没人关注,一旦堵了整栋楼都得停工。Consumer Lag(消费者延迟)就是下水道的流量计,它告诉你生产者往里灌水的速度和消费者抽水的速度差了多少。这个差值持续增大,说明系统出了问题;差值突然归零,也可能出问题(消费者挂了,offset 不再提交,Lag 反而看起来正常)。

这篇文章覆盖 Kafka、RabbitMQ、Redis 三种主流消息队列的监控方案,从指标采集、Prometheus 规则、告警阈值到排障思路,都是我在生产环境踩过坑的实战经验。不讲理论,直接上配置。

为什么 Consumer Lag 是最核心的指标

先说清楚 Lag 到底是什么。Kafka 里每个分区有一个 LogEndOffset(LEO,最新消息位置),消费者组有一个 CurrentOffset(已提交位置)。两者的差值就是 Lag:

Lag = LogEndOffset - CurrentOffset

Lag 为 0 说明消费者完全跟上。但生产环境小幅波动很正常,单分区 Lag < 100 基本不用管。真正要警惕的是三种模式:

Lag 模式典型原因危险程度
持续线性增长消费速度 < 生产速度,处理能力不足高,不处理会雪崩
突然跳变消费者重启/崩溃后 offset 回退,或生产者批量灌数据中,需确认是否预期
突然归零消费者挂了或跳过提交,offset 停滞但 LEO 也没涨低但不正常,常被误判为"健康"

第三种最坑人。我见过一个案例:消费者线程死锁,但 auto.commit.enable=true 还在自动提交 offset,导致消息被标记为"已消费"但实际没处理。Lag 显示 0,业务方以为一切正常,直到下游发现数据丢了三天。

结论:Lag 只看绝对值不够,必须结合 Lag 增长率和消费速率一起看。 一个 Lag=10000 但以每秒 500 的速度下降的队列,远比 Lag=200 但以每秒 10 的速度上涨的队列健康。

Kafka 消费延迟监控全方案

方案一:kafka-exporter(轻量,推荐起步)

社区最常用的方案是 danielqsj/kafka_exporter,部署一个Exporter 容器,自动采集所有 topic 和 consumer group 的 Lag 指标,直接暴露 Prometheus 格式。

# docker-compose.yml
services:
  kafka-exporter:
    image: danielqsj/kafka_exporter:v1.9.0
    command:
      - --kafka.server=kafka1:9092
      - --kafka.server=kafka2:9092
      - --kafka.server=kafka3:9092
      - --topic.filter=^.*$          # 采集所有 topic
      - --group.filter=^.*$          # 采集所有消费组
      - --concurrent.enable          # 并发采集,大集群必开
      - --concurrent.max=200
      - --log.enable-sarama=false   # 关掉 sarama 日志刷屏
    ports:
      - "9308:9308"
    restart: unless-stopped

部署完后访问 http://exporter:9308/metrics,关键指标长这样:

# 当前消费组已提交 offset
kafka_consumergroup_current_offset{consumergroup="order-consumer",topic="orders",partition="0"} 12345678
# 分区最新 offset(LEO)
kafka_topic_partition_leader{topic="orders",partition="0"} 12355678
# 算出来的 Lag(exporter 自动算好)
kafka_consumergroup_lag{consumergroup="order-consumer",topic="orders",partition="0"} 10000

坑提示:kafka_exporter 的 Lag 计算依赖消费者组提交的 offset。如果消费者用的是手动提交且提交间隔很长(比如 60 秒一次),exporter 看到的 Lag 会有数十秒到分钟的延迟。这不是 exporter 的 bug,是 offset 提交机制决定的。如果你的业务对实时性要求高,看后面的方案三。

方案二:JMX Exporter(采集 Broker 内部指标)

kafka_exporter 只看消费端,Broker 自身的健康度需要 JMX 指标。用 jmx_exporter 以 Java Agent 方式挂载到 Kafka 进程:

# kafka 启动参数加一行
KAFKA_OPTS="$KAFKA_OPTS -javaagent:/opt/jmx_exporter/jmx_prometheus_javaagent-0.20.0.jar=9404:/opt/jmx_exporter/kafka.yml"

JMX 采集到的关键指标:

指标含义关注点
kafka_server_replicamanager_underreplicatedpartitions副本同步落后的分区数>0 持续 5 分钟必须告警
kafka_server_requesthandler_avgidlepercent请求线程空闲率<20% 说明 Broker CPU 快撑不住
kafka_network_requestqueue_size请求队列长度持续增长说明处理不过来
kafka_log_log_size每个 topic 的日志大小配合保留策略看磁盘压力

方案三:直接消费 __consumer_offsets(实时性最高)

前面两个方案都有延迟——kafka_exporter 依赖 offset 提交,JMX 看 Broker 状态。如果你要毫秒级的 Lag 监控,最彻底的办法是直接消费 Kafka 内部的 __consumer_offsets topic。

这个 topic 记录了所有消费者组的 offset 提交事件。自己写一个消费者订阅它,解析出每个 group+topic+partition 的最新 offset,再和 LO 对比算出真实处理位置。Flink 2.0 就是这么干的,把 Lag 监控延迟从分钟级压到秒级(参考 Flink 2.0 的 Kafka Source 优化)。

这个方案的代价是开发成本高。普通业务用 kafka_exporter + 30 秒轮询足够了。只有风控、实时推荐这类对延迟极度敏感的场景才需要上方案三。

PromQL 告警规则

三个层次的告警,从轻到重:

# 告警规则文件:kafka-lag-alerts.yml
groups:
  - name: kafka_consumer_lag
    interval: 30s
    rules:
      # P2:单分区 Lag 超过 5000,持续 5 分钟
      - alert: KafkaConsumerLagHigh
        expr: |
          sum by (consumergroup, topic) (
            kafka_consumergroup_lag
          ) > 5000          
        for: 5m
        labels:
          severity: P2
        annotations:
          summary: "Kafka 消费组 {{ $labels.consumergroup }} 的 {{ $labels.topic }} Lag 过高"
          description: "当前总 Lag: {{ $value }},超过阈值 5000"

      # P1:Lag 持续增长 10 分钟未回落
      - alert: KafkaConsumerLagIncreasing
        expr: |
          deriv(
            sum by (consumergroup, topic) (
              kafka_consumergroup_lag
            )[10m:1m]
          ) > 0          
        for: 10m
        labels:
          severity: P1
        annotations:
          summary: "Kafka Lag 持续增长: {{ $labels.consumergroup }} / {{ $labels.topic }}"
          description: "Lag 在过去 10 分钟内持续上升,消费者可能已停止或处理能力不足"

      # P0:Lag 突然归零(疑似消费者挂掉后停止提交)
      - alert: KafkaConsumerLagDroppedToZero
        expr: |
          (sum by (consumergroup, topic) (kafka_consumergroup_lag)) == 0
          and
          (sum by (consumergroup, topic) (kafka_consumergroup_lag offset 5m)) > 10000          
        for: 2m
        labels:
          severity: P0
        annotations:
          summary: "Kafka Lag 突然归零: {{ $labels.consumergroup }}"
          description: "5 分钟前 Lag > 10000,现在归零,疑似消费者崩溃或 offset 异常"

第三条规则是很多人忽略的。Lag 突然归零不一定是好事,很可能是消费者进程挂了、offset 不再更新,而恰好生产者也没新消息。这种情况下 LEO - CurrentOffset = 0,但消息根本没被消费。

分级阈值参考

不同业务的容忍度不同,下面是一个生产环境验证过的阈值表,按业务重要性分级:

业务级别单分区 Lag 警戒总 Lag 警戒增长率警戒告警级别
核心(支付/交易)5002000>50/minP1
重要(订单/通知)200010000>200/minP2
普通(日志/统计)1000050000>500/minP3

RabbitMQ 队列深度监控

Kafka 看的是 offset 差值,RabbitMQ 直接看队列里的消息条数——官方叫 queue depth(队列深度)。含义比 Kafka 的 Lag 更直观:队列里堆了多少没被消费的消息。

内置 Prometheus 插件

RabbitMQ 3.8+ 自带 rabbitmq_prometheus 插件,一行命令启用:

# 在每个 RabbitMQ 节点执行
rabbitmq-plugins enable rabbitmq_prometheus
# 默认暴露在 15692 端口

关键指标:

指标含义告警参考
rabbitmq_queue_messages_ready队列中待消费的消息数>10000 持续 5 分钟
rabbitmq_queue_messages_unacknowledged已投递但未 ack 的消息数>5000 说明消费者处理慢或卡住
rabbitmq_queue_messagesready + unacked 总和综合积压指标
rabbitmq_consumers队列的消费者数量=0 立即告警,消费者全挂了

PromQL 告警规则

groups:
  - name: rabbitmq_queue_depth
    interval: 30s
    rules:
      # P1:消费者全部消失
      - alert: RabbitMQNoConsumers
        expr: rabbitmq_consumers == 0
        for: 1m
        labels:
          severity: P1
        annotations:
          summary: "RabbitMQ 队列 {{ $labels.queue }} 没有消费者"
          description: "队列 {{ $labels.queue }} 在 {{ $labels.vhost }} 下无消费者,消息会一直堆积"

      # P2:队列深度持续过高
      - alert: RabbitMQQueueBacklog
        expr: rabbitmq_queue_messages > 10000
        for: 5m
        labels:
          severity: P2
        annotations:
          summary: "RabbitMQ 队列 {{ $labels.queue }} 积压"
          description: "队列 {{ $labels.queue }} 消息数: {{ $value }}"

      # P1:unacked 消息过多(消费者处理慢或卡住)
      - alert: RabbitMQUnackedHigh
        expr: rabbitmq_queue_messages_unacknowledged > 5000
        for: 3m
        labels:
          severity: P1
        annotations:
          summary: "RabbitMQ unacked 消息过多"
          description: "队列 {{ $labels.queue }} 有 {{ $value }} 条 unacked 消息,消费者可能卡住"

RabbitMQ 和 Kafka 的关键区别:RabbitMQ 是推模式(Push),Broker 主动把消息推给消费者;Kafka 是拉模式(Pull),消费者按需拉取。这意味着 RabbitMQ 的 unacked 指标特别重要——消费者收到消息但迟迟不 ack,说明它在处理过程中卡住了(比如等数据库锁、等下游 HTTP 响应)。这种卡顿不会体现在 ready 消息数上,只体现在 unacked 上。

死信队列监控

RabbitMQ 里消费失败的消息会进死信队列(DLX)。死信队列的消息数也是一个重要信号:

      # P3:死信队列有消息(业务异常,不紧急但需关注)
      - alert: RabbitMQDeadLetterQueue
        expr: |
          rabbitmq_queue_messages{queue=~".*\\.dlq$|.*dead_letter.*"} > 0          
        for: 10m
        labels:
          severity: P3
        annotations:
          summary: "死信队列 {{ $labels.queue }} 有消息"
          description: "死信队列有 {{ $value }} 条消息,检查上游消费失败原因"

Redis 队列监控

很多中小团队用 Redis List 做轻量消息队列(LPUSH 入队,BRPOP 出队)。Redis 本身不是专业消息队列,但胜在简单、延迟低。监控方案和 Kafka/RabbitMQ 不同——Redis 没有内置的"队列 Lag"指标,得自己算。

用 Redis Exporter 采集

# prometheus-redis-exporter
services:
  redis-exporter:
    image: oliver006/redis_exporter:v1.66.0
    command:
      - --redis.addr=redis:6379
      - --redis.password=${REDIS_PASSWORD}
      # 检查指定 key 的 List 长度
      - --check-key-groups=true
      - --check-keys=queue:*,task:*,*:pending
    ports:
      - "9121:9121"

check-keys 参数会让 exporter 定期执行 LLEN 获取 List 的长度。对应的指标:

redis_key_size{key="queue:order_tasks"} 1532
redis_key_size{key="queue:email_tasks"} 0

告警规则

groups:
  - name: redis_queue_depth
    interval: 30s
    rules:
      - alert: RedisQueueBacklog
        expr: redis_key_size{key=~"queue:.*"} > 5000
        for: 5m
        labels:
          severity: P2
        annotations:
          summary: "Redis 队列 {{ $labels.key }} 积压"
          description: "队列长度: {{ $value }}"

      # 延迟队列(Sorted Set)的超时检测
      - alert: RedisDelayedTaskOverdue
        expr: |
          redis_db_keys{db="0"}  # 结合自定义脚本
          # 更好的做法:写个脚本定期 ZCOUNT 统计到期未处理的任务数          
        for: 5m
        labels:
          severity: P2
        annotations:
          summary: "Redis 延迟队列有过期未处理任务"

Redis 队列监控的痛点是 exporter 只能做简单的 LLEN/ZCARD,如果你用了延迟队列(Sorted Set + 时间戳),需要写自定义脚本统计"已到期但未处理"的任务数:

#!/usr/bin/env python3
"""redis_delayed_queue_exporter.py
定期扫描 Redis Sorted Set 类型的延迟队列,
统计已到期但尚未被消费的任务数量。
"""
import time
import redis
from prometheus_client import CollectorRegistry, Gauge, write_to_textfile

r = redis.Redis(host='redis', port=6379, decode_responses=True)
registry = CollectorRegistry()

overdue_tasks = Gauge(
    'redis_delayed_queue_overdue',
    'Number of overdue tasks in delayed queue',
    ['queue_name'],
    registry=registry
)

now = time.time()
# 扫描所有延迟队列 key(约定 key 前缀为 delayed:)
for key in r.scan_iter(match='delayed:*'):
    # ZCOUNT key min max: 统计 score <= now 的元素数
    count = r.zcount(key, 0, now)
    overdue_tasks.labels(queue_name=key).set(count)

write_to_textfile('/var/lib/node_exporter/textfile/redis_delayed.prom', registry)

丢进 crontab 每分钟跑一次,配合 node_exporter 的 textfile collector 就能被 Prometheus 采集。

消费积压的根因排查

监控告警只是第一步。告警响了之后,怎么定位根因?下面是我总结的四象限排查法,按"问题出在生产者还是消费者"和"是突发还是持续"两个维度划分。

四象限定位法

突发(分钟级)持续(小时级)
生产者流量突增(促销/定时任务/消息重放)生产速率长期 > 消费速率
消费者消费者崩溃/重启/OOM/死锁消费逻辑慢(慢查询/下游依赖慢)

突发 + 生产者:最常见。大促、定时跑批、上游补偿性重发。这种积压通常会在流量回落后自行消化。判断标准是看 Lag 增长率——如果增长率开始下降,说明生产速率在回落,等就行。但要把这个判断写进告警,别每次都人工猜。

持续 + 生产者:生产速率长期高于消费速率。可能是新上线了一个数据源,或者上游系统扩容了但消费侧没跟上。这种问题加消费者实例能缓解,但根本解法要么是限流(生产端),要么是扩容(消费端),取决于业务。

突发 + 消费者:消费者挂了。查消费者进程状态、日志、GC 情况。Kafka 场景下,kafka-consumer-groups.sh --describe --group <group> 能看到每个消费者的 host 和分配的 partition。如果消费者列表为空,说明全挂了。

持续 + 消费者:消费逻辑慢。这是最复杂的。先看 jstack 或 pprof 定位消费者线程在干什么——是在等数据库、等 HTTP 响应、还是在做 CPU 密集计算。常见原因是下游依赖变慢(数据库慢查询、第三方接口超时),消费者本身没问题但被下游拖累。

Kafka 排查命令速查

# 1. 查看消费组状态和 Lag
kafka-consumer-groups.sh --bootstrap-server kafka1:9092 \
  --describe --group order-consumer

# 输出示例:
# GROUP           TOPIC       PARTITION  CURRENT-OFFSET  LOG-END-OFFSET  LAG     CONSUMER-ID
# order-consumer  orders      0          12345678        12355678        10000   consumer-1
# order-consumer  orders      1          23456789        23456789        0       consumer-1
# order-consumer  orders      2          34567890        34567990        100     consumer-2
# order-consumer  orders      3          -               45678901        45678901 -  ← 没有消费者!

# 2. 查看消费组成员(确认消费者是否活着)
kafka-consumer-groups.sh --bootstrap-server kafka1:9092 \
  --describe --group order-consumer --members --verbose

# 3. 查看 topic 的分区分布
kafka-topics.sh --bootstrap-server kafka1:9092 \
  --describe --topic orders

# 4. 实时监控 Lag 变化(每 5 秒刷新)
watch -n 5 'kafka-consumer-groups.sh --bootstrap-server kafka1:9092 --describe --group order-consumer | tail -n +3 | sort -k6 -rn'

关键判断:如果某个分区的 CONSUMER-ID 为空(显示 -),说明这个分区没有被分配给任何消费者。原因可能是消费者实例数少于分区数,或者发生了 rebalance。这种 Lag 只能通过增加消费者或修复 rebalance 解决。

RabbitMQ 排查命令

# 查看队列详情(消费者数、消息数、unacked 数)
rabbitmqctl list_queues name messages consumers messages_unacknowledged

# 查看消费者连接
rabbitmqctl list_consumers | grep <queue_name>

# 查看连接详情(排查消费者卡住)
rabbitmqctl list_connections name peer_host state channels

自动化扩容与自愈

监控做好了,下一步是让系统自己处理积压。手动加消费者是救火,自动扩容才是长期解法。

KEDA + Kafka 指标驱动扩缩容

KEDA(Kubernetes Event-Driven Autoscaling)是基于 Kafka Lag 等事件指标自动扩缩容 Deployment 的工具。核心逻辑:Lag 高了自动加 Pod,Lag 降了自动减 Pod。

# ScaledObject: 根据 Kafka Lag 自动扩缩容消费者
apiVersion: keda.sh/v1alpha1
kind: ScaledObject
metadata:
  name: kafka-order-consumer
  namespace: production
spec:
  scaleTargetRef:
    name: order-consumer-deployment
  minReplicaCount: 2        # 最少 2 个副本
  maxReplicaCount: 20       # 最多 20 个副本
  pollingInterval: 30       # 每 30 秒检查一次 Lag
  cooldownPeriod: 300        # 缩容冷却 5 分钟,防抖
  triggers:
    - type: kafka
      metadata:
        bootstrapServers: kafka1:9092,kafka2:9092,kafka3:9092
        consumerGroup: order-consumer
        topic: orders
        lagThreshold: "1000"   # 单分区 Lag 超过 1000 就扩容
        offsetResetPolicy: latest
        partitionLimitation: "0,1,2,3"  # 可选:只监控指定分区

这个配置的效果:当 order-consumer 组在 orders topic 上的 Lag 超过 1000,KEDA 自动增加 Pod 数量。Lag 降回 1000 以下并稳定 5 分钟后,自动缩容。

坑提示:Kafka 的消费者数不能超过分区数。如果你的 topic 只有 6 个分区,Deployment 扩到 20 个 Pod,多余的 14 个 Pod 会空闲等待。扩容前先确认分区数,不够先做分区扩容(kafka-topics.sh --alter --partitions)。

消费者端的反压保护

光扩容不够,消费端还要有自我保护——处理不过来时主动减慢消费,而不是被压垮。Spring Boot Kafka 里用 MaxPollRecords 控制每次拉取的消息数:

# application.yml
spring:
  kafka:
    consumer:
      group-id: order-consumer
      max-poll-records: 100          # 每次最多拉 100 条
      max-poll-interval-ms: 300000   # 两次 poll 间隔上限 5 分钟
      properties:
        max.partition.fetch.bytes: 1048576  # 每分区单次拉取上限 1MB

max-poll-interval-ms 很关键。如果消费者处理一批消息超过这个时间没提交,Kafka 会认为消费者"卡死",触发 rebalance 把分区分给别人。默认值 5 分钟对慢任务可能不够,但设太大又会导致真正卡死的消费者迟迟不被踢出。一般设成你 P99 处理时间的 3 倍。

监控面板设计

Grafana 面板要解决一个问题:一眼看出哪个消费组在积压,而不是在 50 条指标里捞针。

我推荐分三个面板:

面板 1:总览看板(全量消费组 Lag 排名)

# 按 group 汇总总 Lag,降序排列
topk(20, sum by (consumergroup) (kafka_consumergroup_lag))

用 Stat 面板或 Bar Gauge,一眼看出哪个组积压最严重。颜色编码:绿色 <1000,黄色 1000-10000,红色 >10000。

面板 2:单消费组详情(Lag 趋势 + 分区分布)

# Lag 趋势(时间序列图)
sum by (topic) (kafka_consumergroup_lag{consumergroup="$group"})

# 各分区 Lag(表格,高亮最大值)
kafka_consumergroup_lag{consumergroup="$group", topic="$topic"}

面板 3:消费速率(生产速率 vs 消费速率对比)

# 消费速率(每秒消费的消息数)
rate(kafka_consumergroup_current_offset{consumergroup="$group"}[1m])

# 生产速率(每秒新增的消息数)
rate(kafka_topic_partition_leader{topic="$topic"}[1m])

两条线放一张图上。生产速率 > 消费速率时 Lag 必然涨,这是积压的根本信号。

生产环境实践建议

基于踩过的坑,总结几条实战经验:

1. 先建基线再定阈值。 别一上来就抄别人的阈值。先跑一周,统计 Lag 的正常分布(P50、P95、P99),再在 P99 基础上设警戒线。不同业务的 Lag 基线差异巨大——日志消费的 Lag 天然就高,交易消费的 Lag 应该趋近于零。

2. 分区级监控比汇总监控更重要。 总 Lag 正常不代表没问题。6 个分区里 5 个 Lag=0、1 个 Lag=60000,汇总后看起来"还行",但那个分区的消费者可能已经挂了。告警一定要按 consumergroup + topic + partition 维度设,至少按 consumergroup + topic

3. 告警要带上下文。 光告"Lag 10000"没用,值班人不知道是涨还是跌、涨了多久。告警 annotation 里带上当前 Lag、5 分钟前 Lag、增长率、消费者数。值班人一眼能判断是该立刻处理还是等它自己消化。

4. 区分"预期积压"和"异常积压"。 定时任务跑批、数据迁移、消息重放,这些操作天然会产生积压。把这类预期内的高 Lag 告警静默掉,否则告警疲劳会让你忽略真正的故障。

5. 消费者健康检查不只看进程存活。 进程活着不代表在消费。在消费者端打一个 heartbeat 指标——每次成功消费一条消息就更新时间戳。如果 heartbeat 超过 5 分钟没更新,说明消费者虽然活着但实际卡住了。

总结

消息队列监控的核心不是监控工具多花哨,而是你能不能在告警响起的 30 秒内回答三个问题:积压在哪、涨还是跌、消费者活着没。

Kafka 看 Consumer Lag,RabbitMQ 看 ready + unacked,Redis 自己算 List 长度——指标不同,但监控思路一样:先建基线,再设分级阈值,最后配自动化扩容。Lag 归零不一定是好事,Lag 涨了不一定是坏事,关键是看趋势和增长率。

最后提醒一句:监控做得再好,也救不了设计有缺陷的消费逻辑。如果消费者单线程同步调下游慢接口,加再多监控也只是让你更早看到积压。根本解法是异步化、批量化和限流。监控是眼睛,不是手脚。

参考资料与致谢

本文在撰写过程中参考了以下资料,感谢原作者的贡献:

  1. Kafka消息消费卡住了?手把手教你用Offset Explorer监控Lag与排查积压 — CSDN专栏,Consumer Lag 核心指标与积压排查实战
  2. Kafka中replica offset滞后(LAG)过大的常见原因有哪些? — CSDN问答,分区副本 Lag 异常的五维归因分析
  3. SpringBoot消息积压排查:监控与扩容策略 — CSDN博客,消息积压根因与扩容策略
  4. Flink 2.0新特性深度解析:Kafka Lag监控与性能优化实践 — 百度云,Flink 2.0 直接消费 __consumer_offsets 的实时 Lag 监控方案
  5. RQMZ系统中任务队列积压导致延迟,如何优化消费速率? — CSDN问答,四维协同治理模型与动态扩缩容策略