第36讲 | 监控体系搭建
导读:没有监控的集群等于裸奔——本讲梳理 Kafka 核心指标清单,并手把手搭建 Prometheus + Grafana + jmx_exporter 全链路监控大盘。
本讲目标
一、监控什么:四大维度指标清单
1.1 可用性指标(最高优先级,出问题就是事故)
| kafka.controller:type=KafkaController,name=ActiveControllerCount | 当前活跃 controller 数量 | 集群总和恒为 1,≠1 立即严重告警 |
| kafka.server:type=ReplicaManager,name=UnderReplicatedPartitions | ISR 落后于 Leader 的分区数 | > 0 持续 5 分钟告警 |
| kafka.server:type=ReplicaManager,name=OfflinePartitionsCount | 完全不可用(无 Leader)的分区数 | > 0 立即严重告警 |
| kafka.server:type=ReplicaManager,name=UnderMinIsrPartitionCount | ISR 数 < min.insync.replicas 的分区数 | > 0 告警(acks=all 会开始拒绝写入) |
| kafka.controller:type=KafkaController,name=OfflinePartitionsCount(controller 侧) | 控制器视角的离线分区 | 同上 |
UnderReplicatedPartitions 为什么最重要?它等于"副本数 – ISR 数"大于 0 的分区数。副本可能因为 Follower 落后、broker 宕机、磁盘满而被踢出 ISR。持续大于 0 说明复制链路有问题——此时集群仍可用,但冗余度已经下降,再挂一台 broker 就可能出现数据丢失或分区不可用。它是故障的领先指标,而不是滞后指标。
1.2 性能指标(决定用户体验)
| kafka.server:type=BrokerTopicMetrics,name=BytesInPerSec | 每秒写入字节 | 容量水位,超单机承载 70% 需扩容 |
| kafka.server:type=BrokerTopicMetrics,name=BytesOutPerSec | 每秒读出字节 | 消费放大效应 |
| kafka.server:type=BrokerTopicMetrics,name=MessagesInPerSec | 每秒消息数 | 业务侧流量画像 |
| kafka.network:type=RequestMetrics,name=TotalTimeMs,request=Produce | Produce 请求总耗时(分位数) | P99 持续 > 100ms 排查 |
| kafka.network:type=RequestMetrics,name=TotalTimeMs,request=Fetch | Fetch 请求总耗时 | 消费延迟的直接体现 |
| kafka.network:type=RequestChannel,name=RequestQueueSize | 请求队列堆积 | > 100 说明网络线程不够 |
| kafka.network:type=RequestChannel,name=ResponseQueueSize | 响应队列堆积 | 同上,IO 线程瓶颈 |
| kafka.log:type=LogFlushRateAndTimeMs | 刷盘频率与耗时 | 落后磁盘 IO 时激增 |
TotalTimeMs 由四部分构成:RequestQueueTimeMs + LocalTimeMs + RemoteTimeMs + ResponseQueueTimeMs。拆开看就能定位瓶颈在排队(线程不够)、本机处理(磁盘慢)还是等 ISR 副本(远端 broker 慢),这是请求延迟排查的第一把尺子。
1.3 一致性与控制器指标
- kafka.server:type=ReplicaManager,name=IsrShrinksPerSec / IsrExpandsPerSec:ISR 收缩/扩张速率。短暂抖动正常,频繁收缩-扩张循环(flapping)说明网络或 GC 有问题;
- kafka.controller:type=ControllerStats,name=LeaderElectionRateAndTimeMs:Leader 选举频率,非计划内选举频繁要查 broker 掉线原因;
- kafka.controller:type=ControllerStats,name=UncleanLeaderElectionsPerSec:非干净选举,> 0 意味着可能丢数据(unclean.leader.election.enable=false 的集群应为 0);
- KRaft 元数据:kafka.controller:type=KafkaController,name=MetadataErrorCount 与 metadata log 滞后(LastAppliedRecordLagMs)。
1.4 消费者 Lag(业务视角最重要的指标)
Lag = LogEndOffset – ConsumerOffset。Lag 监控有两种途径:
告警按业务定:实时风控场景 Lag > 1000 或分钟级延迟即告警;离线报表场景可以放宽到百万级。
1.5 主机与 JVM 指标
不要只盯 Kafka:CPU iowait、磁盘 util 与 await、网卡包丢弃、PageCache 命中率,以及 JVM 的 GC 暂停时间(jvm.gc.pause,超过 1 秒的 Full GC 会直接触发 session 超时风险)。Kafka 的很多"软件故障"本质是主机故障。
二、监控架构与数据流
#mermaid-svg-ivqm6e1eDn98GRZX{font-family:\”trebuchet ms\”,verdana,arial,sans-serif;font-size:16px;fill:#333;}@keyframes edge-animation-frame{from{stroke-dashoffset:0;}}@keyframes dash{to{stroke-dashoffset:0;}}#mermaid-svg-ivqm6e1eDn98GRZX .edge-animation-slow{stroke-dasharray:9,5!important;stroke-dashoffset:900;animation:dash 50s linear infinite;stroke-linecap:round;}#mermaid-svg-ivqm6e1eDn98GRZX .edge-animation-fast{stroke-dasharray:9,5!important;stroke-dashoffset:900;animation:dash 20s linear infinite;stroke-linecap:round;}#mermaid-svg-ivqm6e1eDn98GRZX .error-icon{fill:#552222;}#mermaid-svg-ivqm6e1eDn98GRZX .error-text{fill:#552222;stroke:#552222;}#mermaid-svg-ivqm6e1eDn98GRZX .edge-thickness-normal{stroke-width:1px;}#mermaid-svg-ivqm6e1eDn98GRZX .edge-thickness-thick{stroke-width:3.5px;}#mermaid-svg-ivqm6e1eDn98GRZX .edge-pattern-solid{stroke-dasharray:0;}#mermaid-svg-ivqm6e1eDn98GRZX .edge-thickness-invisible{stroke-width:0;fill:none;}#mermaid-svg-ivqm6e1eDn98GRZX .edge-pattern-dashed{stroke-dasharray:3;}#mermaid-svg-ivqm6e1eDn98GRZX .edge-pattern-dotted{stroke-dasharray:2;}#mermaid-svg-ivqm6e1eDn98GRZX .marker{fill:#333333;stroke:#333333;}#mermaid-svg-ivqm6e1eDn98GRZX .marker.cross{stroke:#333333;}#mermaid-svg-ivqm6e1eDn98GRZX svg{font-family:\”trebuchet ms\”,verdana,arial,sans-serif;font-size:16px;}#mermaid-svg-ivqm6e1eDn98GRZX p{margin:0;}#mermaid-svg-ivqm6e1eDn98GRZX .label{font-family:\”trebuchet ms\”,verdana,arial,sans-serif;color:#333;}#mermaid-svg-ivqm6e1eDn98GRZX .cluster-label text{fill:#333;}#mermaid-svg-ivqm6e1eDn98GRZX .cluster-label span{color:#333;}#mermaid-svg-ivqm6e1eDn98GRZX .cluster-label span p{background-color:transparent;}#mermaid-svg-ivqm6e1eDn98GRZX .label text,#mermaid-svg-ivqm6e1eDn98GRZX span{fill:#333;color:#333;}#mermaid-svg-ivqm6e1eDn98GRZX .node rect,#mermaid-svg-ivqm6e1eDn98GRZX .node circle,#mermaid-svg-ivqm6e1eDn98GRZX .node ellipse,#mermaid-svg-ivqm6e1eDn98GRZX .node polygon,#mermaid-svg-ivqm6e1eDn98GRZX .node path{fill:#ECECFF;stroke:#9370DB;stroke-width:1px;}#mermaid-svg-ivqm6e1eDn98GRZX .rough-node .label text,#mermaid-svg-ivqm6e1eDn98GRZX .node .label text,#mermaid-svg-ivqm6e1eDn98GRZX .image-shape .label,#mermaid-svg-ivqm6e1eDn98GRZX .icon-shape .label{text-anchor:middle;}#mermaid-svg-ivqm6e1eDn98GRZX .node .katex path{fill:#000;stroke:#000;stroke-width:1px;}#mermaid-svg-ivqm6e1eDn98GRZX .rough-node .label,#mermaid-svg-ivqm6e1eDn98GRZX .node .label,#mermaid-svg-ivqm6e1eDn98GRZX .image-shape .label,#mermaid-svg-ivqm6e1eDn98GRZX .icon-shape .label{text-align:center;}#mermaid-svg-ivqm6e1eDn98GRZX .node.clickable{cursor:pointer;}#mermaid-svg-ivqm6e1eDn98GRZX .root .anchor path{fill:#333333!important;stroke-width:0;stroke:#333333;}#mermaid-svg-ivqm6e1eDn98GRZX .arrowheadPath{fill:#333333;}#mermaid-svg-ivqm6e1eDn98GRZX .edgePath .path{stroke:#333333;stroke-width:2.0px;}#mermaid-svg-ivqm6e1eDn98GRZX .flowchart-link{stroke:#333333;fill:none;}#mermaid-svg-ivqm6e1eDn98GRZX .edgeLabel{background-color:rgba(232,232,232, 0.8);text-align:center;}#mermaid-svg-ivqm6e1eDn98GRZX .edgeLabel p{background-color:rgba(232,232,232, 0.8);}#mermaid-svg-ivqm6e1eDn98GRZX .edgeLabel rect{opacity:0.5;background-color:rgba(232,232,232, 0.8);fill:rgba(232,232,232, 0.8);}#mermaid-svg-ivqm6e1eDn98GRZX .labelBkg{background-color:rgba(232, 232, 232, 0.5);}#mermaid-svg-ivqm6e1eDn98GRZX .cluster rect{fill:#ffffde;stroke:#aaaa33;stroke-width:1px;}#mermaid-svg-ivqm6e1eDn98GRZX .cluster text{fill:#333;}#mermaid-svg-ivqm6e1eDn98GRZX .cluster span{color:#333;}#mermaid-svg-ivqm6e1eDn98GRZX div.mermaidTooltip{position:absolute;text-align:center;max-width:200px;padding:2px;font-family:\”trebuchet ms\”,verdana,arial,sans-serif;font-size:12px;background:hsl(80, 100%, 96.2745098039%);border:1px solid #aaaa33;border-radius:2px;pointer-events:none;z-index:100;}#mermaid-svg-ivqm6e1eDn98GRZX .flowchartTitleText{text-anchor:middle;font-size:18px;fill:#333;}#mermaid-svg-ivqm6e1eDn98GRZX rect.text{fill:none;stroke-width:0;}#mermaid-svg-ivqm6e1eDn98GRZX .icon-shape,#mermaid-svg-ivqm6e1eDn98GRZX .image-shape{background-color:rgba(232,232,232, 0.8);text-align:center;}#mermaid-svg-ivqm6e1eDn98GRZX .icon-shape p,#mermaid-svg-ivqm6e1eDn98GRZX .image-shape p{background-color:rgba(232,232,232, 0.8);padding:2px;}#mermaid-svg-ivqm6e1eDn98GRZX .icon-shape .label rect,#mermaid-svg-ivqm6e1eDn98GRZX .image-shape .label rect{opacity:0.5;background-color:rgba(232,232,232, 0.8);fill:rgba(232,232,232, 0.8);}#mermaid-svg-ivqm6e1eDn98GRZX .label-icon{display:inline-block;height:1em;overflow:visible;vertical-align:-0.125em;}#mermaid-svg-ivqm6e1eDn98GRZX .node .label-icon path{fill:currentColor;stroke:revert;stroke-width:revert;}#mermaid-svg-ivqm6e1eDn98GRZX :root{–mermaid-font-family:\”trebuchet ms\”,verdana,arial,sans-serif;}
Pushgateway / 日志
Broker1 + jmx_exporter :7071
Prometheus
Broker2 + jmx_exporter :7072
Broker3 + jmx_exporter :7073
Grafana 大盘 + 告警
AdminClient Lag 巡检工具
Alertmanager → 电话/企微
为什么选 jmx_exporter + Prometheus 而不是 Jolokia + Zabbix?Kafka 所有指标都通过 JMX 暴露,jmx_exporter 以 Java Agent 方式内嵌在 broker 进程内,零额外进程、无认证需求、Pull 模型天然适合多集群联邦,且 Kafka 社区有现成 Grafana 大盘模板可复用。
三、实战:搭建监控全链路
3.1 Broker 侧接入 jmx_exporter
下载 kafka_prometheus_jmx_exporter(prometheus/jmx_exporter 仓库 release,如 jmx_prometheus_javaagent-0.20.0.jar)到每台 broker 的 /opt/kafka/libs/。
指标规则文件 /opt/kafka/config/jmx-exporter.yaml(只采集关键指标,降低基数):
lowercaseOutputName: true
rules:
# 可用性四件套
– pattern: kafka.server<type=ReplicaManager, name=(UnderReplicatedPartitions|UnderMinIsrPartitionCount|OfflinePartitionsCount)><>Value
name: kafka_replicamanager_$1
type: GAUGE
# 流量
– pattern: kafka.server<type=BrokerTopicMetrics, name=(BytesInPerSec|BytesOutPerSec|MessagesInPerSec), topic=(.+)><>Count
name: kafka_brokertopicmetrics_$1
labels:
topic: "$2"
type: COUNTER
# 请求耗时
– pattern: kafka.network<type=RequestMetrics, name=TotalTimeMs, request=(.+)><>(Count|Mean|99thPercentile)
name: kafka_request_totaltimems_$2
labels:
request: "$1"
type: GAUGE
# 队列
– pattern: kafka.network<type=RequestChannel, name=(RequestQueueSize|ResponseQueueSize)><>Value
name: kafka_requestchannel_$1
type: GAUGE
# JVM GC
– pattern: java.lang<type=GarbageCollector, name=.+><>CollectionTime
name: jvm_gc_collectiontime_total
type: COUNTER
修改 broker 启动环境变量(bin/kafka-server-start.sh 会读取 KAFKA_JVM_PERFORMANCE_OPTS 之外,我们通过 KAFKA_OPTS 注入 agent):
export JMX_PORT=7071
export KAFKA_OPTS="-javaagent:/opt/kafka/libs/jmx_prometheus_javaagent-0.20.0.jar=7071:/opt/kafka/config/jmx-exporter.yaml"
注意:JMX_PORT 与 agent 端口不要冲突,上面示例中 JMX_PORT 仅为满足 Kafka 脚本开启 JMX 的惯例,agent 用同一端口会报地址占用;生产上常见做法是 agent 用 7071,JMX 本地不开远程。若只需 Prometheus 抓取,JMX_PORT 可以不设,仅保留 agent。
重启 broker 后验证:
curl -s http://kafka-1:7071/metrics | grep underreplicated
3.2 Prometheus 抓取配置
prometheus.yml:
global:
scrape_interval: 15s
scrape_configs:
– job_name: "kafka-brokers"
static_configs:
– targets:
– "kafka-1:7071"
– "kafka-2:7071"
– "kafka-3:7071"
labels:
cluster: "prod-kafka"
– job_name: "pushgateway"
static_configs:
– targets: ["pushgateway:9091"]
告警规则 kafka-alerts.yml:
groups:
– name: kafka–critical
rules:
– alert: UnderReplicatedPartitions
expr: sum(kafka_replicamanager_underreplicatedpartitions) by (cluster) > 0
for: 5m
labels:
severity: warning
annotations:
summary: "集群 {{ $labels.cluster }} 存在 ISR 落后分区"
– alert: OfflinePartitions
expr: sum(kafka_replicamanager_offlinepartitionscount) by (cluster) > 0
for: 1m
labels:
severity: critical
annotations:
summary: "集群 {{ $labels.cluster }} 存在无 Leader 分区,立即处理"
– alert: NoActiveController
expr: sum(kafka_controller_activecontrollercount) by (cluster) != 1
for: 1m
labels:
severity: critical
3.3 Grafana 大盘
导入步骤:Grafana → Dashboards → Import,输入社区模板 ID 11962(kafka-exporter 面板需配合 kafka_exporter,若纯 JMX 方案可用 10991 或自建面板),数据源选 Prometheus。
自建面板必备图表(PromQL):
- 集群写入带宽:sum(rate(kafka_brokertopicmetrics_bytesinpersec[5m])) by (topic)
- ISR 收缩速率:sum(rate(kafka_replicamanager_isrshrinks[5m]))(若采集了 IsrShrinksPerSec)
- Produce P99 延迟:kafka_request_totaltimems_99thpercentile{request="Produce"}
- 请求队列:max(kafka_requestchannel_requestqueuesize)
3.4 Lag 监控的三条路
实战案例(Java):AdminClient Lag 巡检工具
这个工具每 30 秒巡检一次所有消费组,发现 Lag 超阈值即打印告警日志(生产上可接企微/钉钉 webhook),适合作为 Prometheus 之外的兜底与双保险。
Maven 依赖:
<dependencies>
<dependency>
<groupId>org.apache.kafka</groupId>
<artifactId>kafka-clients</artifactId>
<version>3.5.1</version>
</dependency>
<dependency>
<groupId>org.slf4j</groupId>
<artifactId>slf4j-simple</artifactId>
<version>2.0.9</version>
</dependency>
</dependencies>
LagPatrolTool.java:
import org.apache.kafka.clients.admin.AdminClient;
import org.apache.kafka.clients.admin.ListOffsetsResult;
import org.apache.kafka.clients.admin.OffsetSpec;
import org.apache.kafka.clients.consumer.OffsetAndMetadata;
import org.apache.kafka.common.TopicPartition;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
import java.util.Properties;
import java.util.concurrent.ExecutionException;
import java.util.stream.Collectors;
/**
* 消费组 Lag 巡检工具:周期性扫描全部消费组,Lag 超阈值输出告警
*/
public class LagPatrolTool {
private static final String BOOTSTRAP = "localhost:9092";
private static final long LAG_WARN_THRESHOLD = 10_000L; // 单分区 Lag 告警阈值
private static final long INTERVAL_MS = 30_000L;
public static void main(String[] args) throws Exception {
Properties props = new Properties();
props.put("bootstrap.servers", BOOTSTRAP);
props.put("request.timeout.ms", 30000);
try (AdminClient admin = AdminClient.create(props)) {
while (true) {
patrolOnce(admin);
Thread.sleep(INTERVAL_MS);
}
}
}
private static void patrolOnce(AdminClient admin)
throws InterruptedException, ExecutionException {
// 1. 列出所有无 DeleteOffSet 提交偏移的消费组,取其已提交 offset
Map<TopicPartition, OffsetAndMetadata> committed = new HashMap<>();
for (String group : admin.listConsumerGroups().all().get()
.stream().map(g -> g.groupId()).collect(Collectors.toList())) {
try {
Map<TopicPartition, OffsetAndMetadata> offsets =
admin.listConsumerGroupOffsets(group).all().get();
// 记录分区归属,便于输出时还原消费组
offsets.keySet().forEach(tp -> committed.put(tp, offsets.get(tp)));
// 同一分区可能被多组消费:这里为简化演示,按组覆盖;
// 生产实现应保存 group -> tp -> offset 三级映射
} catch (Exception e) {
System.err.printf("跳过无位移的消费组 %s: %s%n", group, e.getMessage());
}
}
if (committed.isEmpty()) {
System.out.println("本轮巡检:无已提交偏移");
return;
}
// 2. 查询每个分区的最新末端位移(LogEndOffset)
Map<TopicPartition, OffsetSpec> specs = new HashMap<>();
committed.keySet().forEach(tp -> specs.put(tp, OffsetSpec.latest()));
Map<TopicPartition, ListOffsetsResult.ListOffsetsResultInfo> endOffsets =
admin.listOffsets(specs).all().get();
// 3. 计算 Lag 并告警
long totalLag = 0;
int warnPartitions = 0;
for (Map.Entry<TopicPartition, OffsetAndMetadata> e : committed.entrySet()) {
TopicPartition tp = e.getKey();
long end = endOffsets.get(tp).offset();
long current = e.getValue().offset();
long lag = Math.max(0, end – current);
totalLag += lag;
if (lag > LAG_WARN_THRESHOLD) {
warnPartitions++;
System.out.printf("[WARN] partition=%s lag=%d (current=%d, end=%d)%n",
tp, lag, current, end);
}
}
System.out.printf("本轮巡检完成: partitions=%d, totalLag=%d, warnPartitions=%d%n",
committed.size(), totalLag, warnPartitions);
}
}
运行方式:mvn compile exec:java -Dexec.mainClass=LagPatrolTool。预期输出:
[WARN] partition=order-events-3 lag=52310 (current=102934, end=155244)
本轮巡检完成: partitions=12, totalLag=61022, warnPartitions=1
关键点解读:
- Lag 计算分两步:listConsumerGroupOffsets 拿消费位点、listOffsets(latest) 拿日志末端位点,两数相减即 Lag,与 kafka-consumer-groups.sh –describe 的口径完全一致;
- 工程假设:示例为聚焦主流程用 Map<TopicPartition, OffsetAndMetadata> 保存位移,多消费组共用分区时后查的组会覆盖前者,生产使用需扩展为 Map<String, Map<TopicPartition, Long>> 三级结构(代码注释已标明);
- 与 kafka_exporter 相比,AdminClient 方式无需额外进程、可深度定制告警逻辑(如按业务线分级阈值),代价是要自己维护巡检调度。
踩坑提示 / 生产建议
本讲小结
- 指标四大维度:可用性(ActiveControllerCount、UnderReplicated、Offline、UnderMinIsr)、性能(吞吐、请求延迟、队列)、一致性(ISR flapping、Unclean 选举)、业务(Lag);
- TotalTimeMs = 排队 + 本地处理 + 远端等待 + 响应排队,拆解定位瓶颈;
- jmx_exporter 以 Java Agent 内嵌,Prometheus Pull 抓取,Grafana 模板 10991/7589 快速起盘;
- AdminClient 巡检工具是 Lag 监控的兜底方案,口径与命令行一致。
网硕互联帮助中心



评论前必须登录!
注册