Files
cs-note/hhs/MQ/10-监控与运维/36-MQ-监控指标与告警.md
T

179 lines
7.0 KiB
Markdown
Raw Normal View History

2026-05-24 20:51:06 +08:00
---
tags:
- MQ
- 监控
- Prometheus
- Grafana
create time: 2026-05-24 19:52
---
# MQ 监控指标与告警
## 概述
消息队列作为系统间的"黑盒"中间件,一旦出问题往往影响全局。本文梳理 MQ 监控的核心指标体系,涵盖 Broker、Producer、Consumer 三个维度,并介绍基于 Prometheus + Grafana 的监控告警实践。
## 正文
### 为什么 MQ 需要专门的监控
消息队列不像 Web 服务那样有直观的 HTTP 状态码,它的健康状况藏在各种内部指标里。一个 Topic 的消费延迟可能已经到了几小时,但表面上看 Producer 和 Consumer 都在正常运行——只是 Consumer 跟不上了。等到下游业务发现数据不一致时,往往已经积重难返。
MQ 监控的核心目标就三个字:**看得见**。看得见消息有没有堆积,看得见 Broker 有没有过载,看得见消费者有没有掉队。
### Broker 核心指标
| 指标 | 含义 | 关注点 |
|------|------|--------|
| 消息入队速率 (MessagesIn/s) | 每秒写入的消息数 | 突增可能表示上游流量异常 |
| 消息出队速率 (BytesOut/s) | 每秒读出的字节数 | 与入队速率对比判断消费是否跟得上 |
| 磁盘使用率 | 日志段占用的磁盘空间 | 超过 80% 就该警惕 |
| 网络 IO | 网卡吞吐量 | 高负载时容易成为瓶颈 |
| 连接数 | 当前活跃的客户端连接 | 突增可能是连接泄漏 |
| ISR 数量 | In-Sync Replicas 数量 | ISR 收缩意味着有 Follower 掉队 |
> [!question]
> ISR(In-Sync Replicas)数量减少时,消息的可靠性会受到什么影响?如果你设置了 `acks=all`,ISR 缩减到 1 会发生什么?
### Producer 核心指标
Producer 侧最需要关注的是**发送成功率**和**发送延迟**。
- **发送成功率**:失败的发送请求占比。如果持续有失败,说明 Broker 端有问题(磁盘满、网络分区)。
- **发送延迟 P99**:99 分位的发送耗时。正常情况下应该在毫秒级,如果飙升到秒级,说明 Broker 压力过大或者网络抖动。
- **重试率**:发送失败后重试的比例。高重试率意味着消息可能乱序(Kafka 中同一 Partition 内的消息顺序靠 offset 保证,但重试可能导致后发的消息先到)。
- **批大小(Batch Size)**:Producer 端的批量发送大小。批越大吞吐越高,但延迟也越大。
### Consumer 核心指标
Consumer 侧最核心的指标是 **Consumer Lag**——消费者当前的消费位置与最新消息之间的差距。
Lag 本质上衡量的是"消费者落后了多远"。如果 Lag 持续增长,说明消费者处理不过来,消息正在积压。
- **消费速率**:每秒处理的消息数,应与 Producer 的入队速率基本持平。
- **Rebalance 次数**:Consumer Group 发生 Rebalance 的频率。频繁 Rebalance 会导致消费暂停,通常是 Consumer 不稳定(频繁重启、处理超时)引起的。
- **处理耗时**:单条消息从业务处理的平均耗时。如果处理耗时接近 `max.poll.interval.ms`,就有触发 Rebalance 的风险。
### Prometheus + Grafana 监控搭建
监控架构分四层:
```mermaid
graph LR
MQ["MQ Cluster"] -->|"指标暴露"| Exporter["MQ Exporter"]
Exporter -->|"pull /metrics"| Prometheus["Prometheus"]
Prometheus -->|"查询"| Grafana["Grafana"]
Prometheus -->|"告警规则"| AlertManager["AlertManager"]
AlertManager -->|"通知"| Notify["DingTalk / PagerDuty"]
```
**Exporter 配置**:Kafka 可以用 `kafka_exporter` 或 JMX Exporter 暴露指标。`kafka_exporter` 轻量级,适合快速接入;JMX Exporter 功能更全,但配置更复杂。
```yaml
# Prometheus 配置示例
scrape_configs:
- job_name: "kafka"
static_configs:
- targets: ["kafka-exporter:9308"]
scrape_interval: 15s
```
**核心 Dashboard 设计**:一个实用的 MQ Dashboard 通常包含以下几个面板:
1. **总览面板**:集群消息入队/出队速率、总 Lag 数、Broker 存活数。
2. **Broker 面板**:每个 Broker 的磁盘使用率、网络 IO、ISR 数量、请求队列深度。
3. **Topic 面板**:每个 Topic 的消息速率、Partition 分布、Lag 趋势。
4. **Consumer Group 面板**:每个 Group 的 Lag、消费速率、Rebalance 次数。
### 告警策略
好的告警策略不是"什么都报",而是"报了就要行动"。三种常用的告警模型:
**基于阈值**:最直观,比如 `Consumer Lag > 10000` 触发告警。适合有明确 SLO 的场景。
**基于趋势**:Lag 的绝对值不重要,重要的是它是否在持续增长。比如"Lag 连续 10 分钟单调递增"就该告警,即使当前 Lag 只有 100。
**基于异常检测**:用统计方法检测指标是否偏离正常范围。比如消费速率突然降到 0,但 Producer 侧没有任何变化,这很可能是 Consumer 挂了。
> [!question]
> 你能设计一个既不会频繁误报、又不会漏报关键问题的 Consumer Lag 告警规则吗?阈值设多少合适?
### Go 代码示例:自定义 Consumer Lag 监控
```go
package main
import (
"context"
"log"
"net/http"
"time"
"github.com/IBM/sarama"
"github.com/prometheus/client_golang/prometheus"
"github.com/prometheus/client_golang/prometheus/promhttp"
)
var (
// consumerLag 每个 Partition 的消费延迟
consumerLag = prometheus.NewGaugeVec(
prometheus.GaugeOpts{
Name: "myapp_consumer_lag",
Help: "Consumer lag per partition",
},
[]string{"topic", "partition", "group"},
)
)
func init() {
prometheus.MustRegister(consumerLag)
}
// collectLag 定期采集 Consumer Lag 并暴露给 Prometheus
func collectLag(brokers []string, group, topic string) {
config := sarama.NewConfig()
client, err := sarama.NewClient(brokers, config)
if err != nil {
log.Fatalf("create client: %v", err)
}
defer client.Close()
ticker := time.NewTicker(10 * time.Second)
for range ticker.C {
partitions, _ := client.Partitions(topic)
for _, p := range partitions {
// 获取最新 offset
newest, _ := client.GetOffset(topic, p, sarama.OffsetNewest)
// 获取消费者组当前 offset(实际需通过 Admin API)
// 这里简化为从外部获取
groupOffset := getGroupOffset(group, topic, p)
lag := float64(newest - groupOffset)
consumerLag.WithLabelValues(topic,
string(rune(p+'0')), group).Set(lag)
}
}
}
func getGroupOffset(group, topic string, partition int32) int64 {
// 实际实现中通过 Kafka Admin API 获取
return 0
}
func main() {
go collectLag([]string{"localhost:9092"}, "my-group", "orders")
http.Handle("/metrics", promhttp.Handler())
log.Println("metrics server on :2112")
log.Fatal(http.ListenAndServe(":2112", nil))
}
```
这段代码通过 Sarama 客户端定期查询每个 Partition 的最新 offset 和消费者组当前 offset,计算差值作为 Lag,然后通过 Prometheus Gauge 暴露。Grafana 中可以直接用 `myapp_consumer_lag` 指标配置 Dashboard 和告警规则。
## 关联笔记
- [[30-MQ-核心概念与选型]]
- [[33-MQ-Kafka-架构与核心机制]]
- [[37-MQ-消费积压治理]]
- [[35-MQ-与-CDC]]