概述
凌晨三点,手机被告警轰炸。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 警戒 | 增长率警戒 | 告警级别 |
|---|---|---|---|---|
| 核心(支付/交易) | 500 | 2000 | >50/min | P1 |
| 重要(订单/通知) | 2000 | 10000 | >200/min | P2 |
| 普通(日志/统计) | 10000 | 50000 | >500/min | P3 |
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_messages | ready + 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 涨了不一定是坏事,关键是看趋势和增长率。
最后提醒一句:监控做得再好,也救不了设计有缺陷的消费逻辑。如果消费者单线程同步调下游慢接口,加再多监控也只是让你更早看到积压。根本解法是异步化、批量化和限流。监控是眼睛,不是手脚。
参考资料与致谢
本文在撰写过程中参考了以下资料,感谢原作者的贡献:
- Kafka消息消费卡住了?手把手教你用Offset Explorer监控Lag与排查积压 — CSDN专栏,Consumer Lag 核心指标与积压排查实战
- Kafka中replica offset滞后(LAG)过大的常见原因有哪些? — CSDN问答,分区副本 Lag 异常的五维归因分析
- SpringBoot消息积压排查:监控与扩容策略 — CSDN博客,消息积压根因与扩容策略
- Flink 2.0新特性深度解析:Kafka Lag监控与性能优化实践 — 百度云,Flink 2.0 直接消费
__consumer_offsets的实时 Lag 监控方案 - RQMZ系统中任务队列积压导致延迟,如何优化消费速率? — CSDN问答,四维协同治理模型与动态扩缩容策略