主题
第 19 章 可观测性与运维
一句话总结本章:「Kafka 不是黑盒,但你不主动看,它就是黑盒」。
学完你会:
- 能背出 Broker / Producer / Consumer 三大类核心 JMX 指标的健康阈值
- 会搭 JMX → JMX Exporter → Prometheus → Grafana 的完整观测链路
- 看
server.log/controller.log/state-change.log/log-cleaner.log不再眼瞎- 熟练操作分区重分配、Leader 优先选举、节点退役 / 上线
- 会做 JVM、OS、磁盘、网络的调优清单
📌 生活类比:可观测性就像「医院体检」。光会看病人脸色(错误日志)远远不够,你要知道:
- 体温 =
UnderReplicatedPartitions(高烧 = 副本掉队)- 心率 =
MessagesInPerSec(异常波动 = 业务异常)- 血压 =
RequestHandlerAvgIdlePercent(持续偏低 = Broker 处理不过来)- X 光 / CT =
kafka-dump-log.sh(看 Segment 内部到底发生了什么)
0. 一张图看清「观测三件套」
┌───────────────────────────────────┐
│ Grafana Dashboard │
│ (UnderReplicated / Lag / TPS …) │
└────────────────┬──────────────────┘
│ PromQL 查询
▼
┌───────────────────────────────────┐
│ Prometheus │
│ (TSDB + Alertmanager 告警) │
└────────────────┬──────────────────┘
│ HTTP scrape (15s)
┌────────────────────────────────┼────────────────────────────────┐
▼ ▼ ▼
┌─────────────────┐ ┌─────────────────┐ ┌─────────────────┐
│ JMX Exporter │ │ JMX Exporter │ │ kafka-exporter │
│ (sidecar :7071) │ │ (sidecar :7071) │ │ (LAG 专用) │
└────────┬────────┘ └────────┬────────┘ └────────┬────────┘
│ JMX │ JMX │ AdminClient
▼ ▼ ▼
Broker 1 Broker 2 连任意 Broker- JMX Exporter:跑在 Broker 旁边(Java Agent 或 sidecar),把 JMX MBean 转成 Prometheus 文本格式
- kafka-exporter:用 AdminClient 抓 Topic / Group 的 Lag、ISR 等业务指标,Lag 看它最方便
- Prometheus:定时拉取并落盘
- Grafana:展示 + 告警面板
1. 核心 JMX 指标清单(必背)
所有指标都来自 Kafka Broker 的 JMX,路径形如
kafka.server:type=ReplicaManager,name=UnderReplicatedPartitions。
1.1 Broker 健康类(5 大致命指标)
| 指标 | 含义 | 健康阈值 | 异常排查方向 |
|---|---|---|---|
UnderReplicatedPartitions | 没有完整 ISR 的分区数 | = 0 | Follower 掉队 / 网络延迟 / 磁盘慢;查 replica.lag.time.max.ms 与 ReplicaFetcher Lag |
OfflinePartitionsCount | 没有 Leader 的分区数(最致命) | = 0 | Broker 全部宕机或 Controller 异常;立刻看 controller.log |
ActiveControllerCount | 当前 Controller 数 | 整集群恰为 1 | =0:没人当 Controller,集群脑死;>1:脑裂,立刻人工介入 |
LeaderElectionRateAndTimeMs | Leader 选举速率与耗时 | 平时 ≈ 0 | 突增 = Broker 抖动 / 网络问题,看 state-change.log |
UncleanLeaderElectionsPerSec | 不干净选举速率(会丢数据) | = 0 | 已发生丢数据;检查是否误开 unclean.leader.election.enable=true |
1.2 流量与性能类
| 指标 | 含义 | 健康阈值 | 异常排查方向 |
|---|---|---|---|
MessagesInPerSec | 每秒进入消息数(每 Topic / 全局) | 看业务基线 | 突降 = 上游异常;突增 = 上游打满 |
BytesInPerSec | 每秒入流量 | < 网卡带宽 70% | 接近上限要扩容 |
BytesOutPerSec | 每秒出流量 | < 网卡带宽 70% | 同上;通常 OutPerSec 是 InPerSec × 副本数 + Consumer fanout |
RequestHandlerAvgIdlePercent | I/O 线程空闲率 | > 30% | 长期 < 20% 表示 num.io.threads 不够 |
NetworkProcessorAvgIdlePercent | 网络线程空闲率 | > 30% | 长期 < 20% 表示 num.network.threads 不够 |
TotalTimeMs | 单个请求各阶段耗时(百分位) | p99 < 100ms | 拆 5 个阶段看:本地、远端、副本、响应队列、响应发送 |
RequestQueueSize | 请求队列深度 | < queued.max.requests 的 50% | 满了 = 处理不过来 |
LogFlushRateAndTimeMs | 刷盘频率与耗时 | 看磁盘类型 | 持续高耗时 = 磁盘瓶颈 |
1.3 副本与 ISR 类
| 指标 | 含义 | 健康阈值 | 异常排查方向 |
|---|---|---|---|
ISRShrinksPerSec | ISR 收缩速率 | ≈ 0 | Follower 掉队 / 网络抖动 |
ISRExpandsPerSec | ISR 扩张速率 | ≈ 0 | 与 Shrink 配对出现是「ISR 抖动」 |
ReplicaFetcher MaxLag | 副本同步最大滞后 | 几千以内 | 持续高 = 该 Follower 跟不上,检查带宽 / GC |
LeaderCount | Leader 数 | 各 Broker 均匀(差异 < 10%) | 不均匀 → 跑 kafka-leader-election.sh --election-type preferred |
PartitionCount | 分区数 | 各 Broker 均匀 | 不均匀 → kafka-reassign-partitions.sh |
1.4 客户端 Producer 指标
| 指标 | 含义 | 阈值 | 排查方向 |
|---|---|---|---|
record-error-rate | 发送错误速率 | = 0 | 看错误类型(NotLeader / NotEnoughReplicas / TimedOut) |
record-send-rate | 发送成功速率 | 看业务基线 | 突降 = 业务异常或 Broker 拒收 |
record-retry-rate | 重试速率 | 接近 0 | 高 = Leader 切换 / 副本不足 |
request-latency-avg | 平均请求延迟 | < 50ms | p99 单独看 |
request-latency-p99 | p99 延迟 | < 200ms | 高 = 副本同步慢、网络抖动 |
record-queue-time-max | 在 Producer 内排队时间 | < linger.ms × 2 | 高 = batch 满了发不出去 |
bufferpool-wait-time-total | 缓冲池等待时间 | ≈ 0 | 高 = buffer.memory 不够 |
1.5 客户端 Consumer 指标
| 指标 | 含义 | 阈值 | 排查方向 |
|---|---|---|---|
records-lag-max | 最大 Lag(消息条数) | 业务可接受范围 | 高 = 消费跟不上,扩容消费者 / 优化处理 |
records-consumed-rate | 消费速率 | 看业务基线 | 突降 = Rebalance / 业务卡住 |
fetch-latency-avg | 拉取延迟 | < 100ms | 高 = 网络 / Broker 慢 |
fetch-rate | 拉取频率 | 合理 | 过高 = 每次拉得太少;过低 = poll 间隔太长 |
commit-latency-avg | 提交 Offset 延迟 | < 50ms | 高 = __consumer_offsets 副本异常 |
assigned-partitions | 当前分配到的分区数 | > 0 | =0 = 被踢出消费组 |
rebalance-rate-per-hour | 每小时 Rebalance 次数 | < 1 | 高 = max.poll.interval.ms 配置不当 |
⚠️ Lag 的两种语义:
- records-lag:还有多少条没消费(业务关心)
- fetcher-lag:内部 Fetcher 与 Leader 的滞后(运维关心)
2. JMX → Prometheus → Grafana 全链路
2.1 启动 JMX Exporter
JMX Exporter 以 Java Agent 形式注入 Kafka 进程,监听 7071 端口输出 Prometheus 文本:
bash
# 启动 Kafka 时加参数
export KAFKA_OPTS="-javaagent:/opt/jmx_exporter/jmx_prometheus_javaagent.jar=7071:/opt/jmx_exporter/kafka.yml"
bin/kafka-server-start.sh config/server.propertieskafka.yml 是规则文件,决定哪些 MBean 暴露成什么名字(见 code/jmx_exporter_config.yml)。
2.2 Prometheus 抓取
yaml
# prometheus.yml
scrape_configs:
- job_name: kafka-broker
scrape_interval: 15s
static_configs:
- targets: ['kafka1:7071','kafka2:7071','kafka3:7071']
labels:
cluster: prod-1
- job_name: kafka-exporter
static_configs:
- targets: ['kafka-exporter:9308']2.3 Grafana Dashboard 推荐 ID
| ID | 名称 | 说明 |
|---|---|---|
| 7589 | Kafka Exporter Overview | Topic / Group / Lag 总览 |
| 11962 | Kafka (JMX Exporter) | Broker 全维度 |
| 21078 | Kafka Cluster Overview | Confluent 风格 |
| 18276 | Strimzi Kafka | K8s 部署专用 |
📌 直接在 Grafana 「Import dashboard」输入 ID 即可拉取。
2.4 关键告警规则示例
yaml
groups:
- name: kafka.rules
rules:
- alert: KafkaUnderReplicated
expr: sum(kafka_server_replicamanager_underreplicatedpartitions) by (instance) > 0
for: 5m
labels: { severity: warning }
annotations:
summary: "Broker {{ $labels.instance }} 有 {{ $value }} 个分区副本不全"
- alert: KafkaOfflinePartitions
expr: sum(kafka_controller_kafkacontroller_offlinepartitionscount) > 0
for: 1m
labels: { severity: critical }
annotations:
summary: "Kafka 集群有 Offline 分区!"
- alert: KafkaConsumerLagHigh
expr: kafka_consumergroup_lag > 10000
for: 10m
labels: { severity: warning }
annotations:
summary: "Group {{ $labels.consumergroup }} Lag = {{ $value }}"3. 日志定位手册
Kafka 在 logs/ 目录下产出 4 个最重要的日志:
3.1 server.log(80% 问题在这看)
- Broker 启动 / 关闭、客户端连接、请求异常、分区状态变更
- 常见关键词:
ERROR/WARN/Connection reset/Disconnected - 必看:启动时是否成功加入集群(
Registered broker xxx),是否有大量Failed to send错误
3.2 controller.log(控制面专属)
- Controller 选举、分区 Leader 选举、副本分配
- 当
OfflinePartitionsCount报警时,第一时间打开它 - 常见关键词:
elected as the controller/LeaderAndIsr request/failed controller event
3.3 state-change.log(分区状态变化)
- 每个分区从
OfflinePartition → OnlinePartition、Leader 切换、ISR 变更的时间线 - 适合「事后回放」:复盘某个分区为什么 12:34 切了 Leader
3.4 log-cleaner.log(Compaction 专属)
- Log Compaction 线程的工作日志
- 常见关键词:
Cleaner thread/compacted in/aborted - Compaction 不生效时必看:是否报
OutOfMemoryError、是否被某个 Segment 卡死
📌 小技巧:用
tail -F配合grep实时盯关键词:bashtail -F server.log | grep -E 'ERROR|WARN|UnderReplicated'
4. 常用运维命令大全
4.1 分区重分配(最常用,也最容易出事)
场景:新加了 Broker,需要把分区均衡过去;或某 Broker 要下线。
三步法:
bash
# 第 1 步:生成方案 JSON
cat > topics.json <<EOF
{"topics":[{"topic":"learn.18.users"}],"version":1}
EOF
bin/kafka-reassign-partitions.sh --bootstrap-server localhost:9092 \
--topics-to-move-json-file topics.json \
--broker-list "1,2,3,4" \
--generate
# 输出两段 JSON:「Current partition replica assignment」「Proposed partition reassignment configuration」
# 把 Proposed 那段保存为 reassign.json
# 第 2 步:执行(带限速!)
bin/kafka-reassign-partitions.sh --bootstrap-server localhost:9092 \
--reassignment-json-file reassign.json \
--execute \
--throttle 50000000 # 50 MB/s,避免吃光网络
# 第 3 步:验证 + 取消限速
bin/kafka-reassign-partitions.sh --bootstrap-server localhost:9092 \
--reassignment-json-file reassign.json \
--verify⚠️ 务必限速!默认无限速时,重分配会吃光带宽,把业务一起拖垮。 ⚠️ 完成后记得 verify,verify 会清理
__reassign_partitions节点和限速配置,否则限速会永久生效。
4.2 Leader 优先副本选举(修不均衡)
bash
# 给所有 Topic 做 preferred 选举
bin/kafka-leader-election.sh --bootstrap-server localhost:9092 \
--election-type preferred --all-topic-partitions
# 给特定分区做
bin/kafka-leader-election.sh --bootstrap-server localhost:9092 \
--election-type preferred --topic learn.18.users --partition 3📌 「Preferred Leader」= Topic 创建时 replica 列表里第一个副本。集群运行久了,Leader 会因为 Broker 重启等原因「漂移」,跑这个命令可以让 Leader 回到最初的均衡分布。
4.3 节点退役流程
1. 把该 Broker 上所有 Leader 切走(preferred election 或重分配)
2. 用 kafka-reassign-partitions.sh 把所有分区从该 Broker 移走
3. 等所有 Topic 的 ISR 不再包含该 Broker
4. 安全 stop kafka-server-stop.sh
5. (KRaft 模式)用 kafka-storage.sh 注销节点4.4 节点上线
1. 配 broker.id(KRaft 用 node.id)
2. 启动 → 自动加入集群(看 controller.log 是否有 Registered broker)
3. 此时新 Broker 没分区,需要主动重分配让它接活4.5 查看 / 重置消费组 Offset
bash
# 查看所有组
bin/kafka-consumer-groups.sh --bootstrap-server localhost:9092 --list
# 查看具体组的 Lag
bin/kafka-consumer-groups.sh --bootstrap-server localhost:9092 \
--describe --group my-group
# 重置到最早
bin/kafka-consumer-groups.sh --bootstrap-server localhost:9092 \
--group my-group --topic my-topic --reset-offsets --to-earliest --execute
# 重置到指定时间戳
bin/kafka-consumer-groups.sh --bootstrap-server localhost:9092 \
--group my-group --topic my-topic --reset-offsets \
--to-datetime 2026-04-15T00:00:00.000 --execute⚠️ 重置 Offset 必须在消费者下线时操作,否则会失败。
5. JVM 调优要点
5.1 堆大小:不要太大!
bash
# 推荐:6~8 GB
export KAFKA_HEAP_OPTS="-Xms6g -Xmx6g"反直觉的真理:Kafka 不靠 JVM 堆缓存数据,它依赖 Linux PageCache。给 JVM 堆 50 GB 反而害死自己:
- 堆大 → GC 暂停长 → ISR 抖动
- 堆大 → 留给 PageCache 的物理内存少 → 读写性能下降
经验值:64 GB 物理内存的机器,JVM 给 8 GB,剩下 56 GB 全留给 PageCache。
5.2 GC:必选 G1
bash
export KAFKA_JVM_PERFORMANCE_OPTS="-server \
-XX:+UseG1GC \
-XX:MaxGCPauseMillis=20 \
-XX:InitiatingHeapOccupancyPercent=35 \
-XX:+ExplicitGCInvokesConcurrent \
-Djava.awt.headless=true"MaxGCPauseMillis=20:每次 GC 不超过 20 ms(实际可能不达标,但作为目标值)InitiatingHeapOccupancyPercent=35:堆 35% 满就启动并发 GC,避免堆满 STW
5.3 网络与 I/O 线程数
server.properties:
properties
num.network.threads=8 # 默认 3,机器核数多就调高
num.io.threads=16 # 默认 8,I/O 密集型场景翻倍
queued.max.requests=1000 # 请求队列深度
socket.send.buffer.bytes=102400
socket.receive.buffer.bytes=102400
socket.request.max.bytes=104857600 # 100MB,大消息场景调大调参依据:
RequestHandlerAvgIdlePercent长期 < 20% → 加num.io.threadsNetworkProcessorAvgIdlePercent长期 < 20% → 加num.network.threads- 别盲目加,先看指标
6. OS 层调优清单
6.1 内核参数
bash
# /etc/sysctl.conf
# 几乎不用 swap,避免 PageCache 被换出
vm.swappiness = 1
# Dirty Page 比例
vm.dirty_background_ratio = 5 # 5% 开始后台刷
vm.dirty_ratio = 80 # 80% 强制阻塞写
# 网络缓冲
net.core.rmem_max = 67108864
net.core.wmem_max = 67108864
net.ipv4.tcp_rmem = 4096 87380 67108864
net.ipv4.tcp_wmem = 4096 65536 67108864
# 大量连接场景
net.core.somaxconn = 4096
net.ipv4.tcp_max_syn_backlog = 4096执行:sysctl -p
6.2 文件句柄数
bash
# /etc/security/limits.conf
kafka soft nofile 100000
kafka hard nofile 100000📌 为什么要这么大:每个 Topic-Partition 至少占用 3 个 fd(
.log/.index/.timeindex),分区多时轻松上万。
6.3 文件系统
| 文件系统 | 评价 |
|---|---|
| XFS | 首选,大文件顺序写性能最好 |
| ext4 | 可用,性能略差 |
| ZFS | 不推荐(双层 cache 与 PageCache 冲突) |
| 网络盘 (NFS / EBS gp2) | 禁用 用于数据目录,延迟太大 |
挂载参数:noatime,nodiratime(不更新访问时间,省 IO)
6.4 磁盘
- 机械盘:可以用,但单盘吞吐 100 MB/s 左右
- SSD:推荐 NVMe,单盘吞吐 1-3 GB/s
- RAID:用 RAID 10 或多盘
log.dirs多目录,不要 RAID 5(写放大严重)
7. 监控告警最佳实践
7.1 告警分级
| 级别 | 触发条件 | 响应 |
|---|---|---|
| P0 致命 | OfflinePartitions > 0 / ActiveController != 1 | 5 分钟内有人响应 |
| P1 严重 | UnderReplicated > 0 持续 5 分钟 / Lag > 阈值 | 30 分钟内响应 |
| P2 警告 | RequestHandlerIdle < 20% / GC 时间高 | 1 小时内查看 |
| P3 通知 | Topic 增长率异常 / 磁盘 60% | 上班时间处理 |
7.2 三个「上线必须验证」的指标
UnderReplicatedPartitions = 0:副本健康ActiveControllerCount = 1(全集群):Controller 健康UncleanLeaderElectionsPerSec = 0:无数据丢失
7.3 一定要看的「业务指标」
- 各 Topic 的 实际吞吐(消息数 / 字节数)
- 各 Consumer Group 的 Lag 趋势
- 各 Topic 的 平均消息大小(突增意味着业务异常或日志爆炸)
- 生产者错误率(按错误类型分类)
8. 故障演练(运维成熟度的标尺)
每季度做一次:
| 演练 | 操作 | 期望结果 |
|---|---|---|
| Broker 宕机 | kill -9 一个 Broker | Leader 在 30 秒内切走,业务无感 |
| 磁盘满 | dd 把 log 目录填满 | 该 Broker 拒收写入但不 panic,迁出后恢复 |
| 网络分区 | iptables drop 与某 Broker 的连接 | ISR 收缩,业务降级但不崩溃 |
| Controller 宕机 | kill 当前 Controller 节点 | 选举新 Controller,期间元数据变更阻塞 < 30 秒 |
| Schema Registry 挂 | stop SR | Producer/Consumer 走本地缓存,新 Schema 注册失败但已有业务正常 |
9. 横向对比:与其他 MQ 的运维差异
| 维度 | Kafka | RabbitMQ | RocketMQ | Pulsar |
|---|---|---|---|---|
| 元数据 | KRaft / ZK | Mnesia / 内存 | NameServer | ZK + BookKeeper |
| 主要指标 | UnderReplicated / Lag | queue length / consumers | broker offset / consume lag | bookie ledger |
| 重分配 | reassign-partitions(手动) | 队列绑定改路由 | 调整队列分布 | bundle 自动 reassign |
| 重启代价 | 较高(PageCache 重建) | 中(持久化队列恢复慢) | 中 | 低(Broker 无状态) |
| 单机扩展 | 加 broker + reassign | 加 node + mirror | 加 broker + topic 扩展 | 加 broker / 加 bookie 解耦 |
10. 本章面试高频题
Q1:判断 Kafka 集群健康,最先看哪几个指标?
答案(背 5 个):
OfflinePartitionsCount—— 必须 = 0(否则有数据完全无法读写)UnderReplicatedPartitions—— 应 = 0(持续大于 0 说明副本掉队)ActiveControllerCount—— 整集群恰为 1(>1 = 脑裂,=0 = 集群瘫痪)UncleanLeaderElectionsPerSec—— = 0(>0 说明发生了「丢数据」式选举)RequestHandlerAvgIdlePercent—— > 30%(< 20% 说明 I/O 线程不够)
加分:能解释每个指标的物理含义、阈值由来、常见根因。
Q2:Consumer Lag 一直在涨,怎么排查?
结构化答:
- 先确认 Lag 真涨:
kafka-consumer-groups.sh --describe看具体分区的 Lag - 定位是上游写多 / 下游消费慢:对比
MessagesInPerSec与 Consumerrecords-consumed-rate - 如是消费慢:
- 业务处理慢 → 看应用日志、数据库慢查询、外部依赖
- 单线程消费 → 用
max.poll.records+ 多线程处理 - Rebalance 频繁 → 调
max.poll.interval.ms/session.timeout.ms
- 结构性原因:分区数 < 消费者期望并发数 → 加分区
- 临时止血:扩容消费者实例(不超过分区数)
加分项:提到 Lag 也可能是「Heartbeat 正常但消费卡住」;提到要看 time-lag(按时间)而非只看 records-lag(按条数)。
Q3:JVM 堆给多大合适?为什么不能给 50G?
答案:
- 推荐 6~8 GB
- Kafka 的数据缓存依赖 Linux PageCache,不在 JVM 堆里
- 堆大 → GC 暂停长(>1 秒)→ ISR 收缩 / Heartbeat 超时 → Rebalance / Leader 切换
- 堆大 → PageCache 可用内存少 → 读写都需要走磁盘 → 性能急剧下降
加分:能说出 Kafka 用 sendfile 零拷贝、PageCache 命中率是核心性能指标。
Q4:UnderReplicatedPartitions > 0 时怎么处理?
答案分阶段:
- 看
controller.log/server.log:定位是哪个 Broker / Topic / Partition - 看 ReplicaFetcher Lag:是不是某个 Follower 跟不上
- 常见原因:
- Broker 宕机 / 网络隔离 → 重启或修网络
- Broker 磁盘满 → 清日志 / 加盘 / 迁移分区
- GC 暂停长 → 调 JVM
- 副本同步限速 (
replica.fetch.max.bytes/num.replica.fetchers) 配置过低 → 调高
- 临时缓解:把该 Broker 上的 Leader 切走(preferred election 排除该节点)
绝对禁忌:不要打开 unclean.leader.election.enable=true(会丢数据)。
Q5:分区重分配为什么必须 --throttle?
答案:
- 重分配本质是 Follower 从源 Broker 拉数据,会额外占用网络与磁盘 IO
- 不限速时,重分配流量可能吃光带宽,导致:
- 业务 Producer 写入超时
- Consumer 拉数据延迟暴增
- 其他正常副本同步落后 → ISR 抖动 → 雪崩
--throttle给 Leader → Follower 限速,限制重分配的副本同步速率- 完成后必须
--verify,否则限速配置会永久残留
加分:知道 throttle 实际通过 leader.replication.throttled.rate / follower.replication.throttled.rate 配置存到 ZK / KRaft。
Q6:日志想本地保留 3 天 + 远程归档,怎么做?
答案:
- Topic 配置:
retention.ms = 259200000 # 3 天 cleanup.policy = delete - 远程归档方案:
- KIP-405 Tiered Storage(Kafka 3.6+,把老 Segment 异步上传 S3 / OSS)
- 用 Kafka Connect S3 Sink 实时备份
- 用 MirrorMaker 2 复制到归档集群
- 监控:
LogSize是否在按预期下降- 远程存储的对象数 / 体积是否持续增长
11. 小结
- 可观测性 = 5 大致命指标 + 3 类客户端指标 + 4 个日志文件
- JMX → Prometheus → Grafana 是事实标准链路;用 kafka-exporter 补 Lag 指标
- JVM 调优黄金法则:堆给 6-8G、剩下让给 PageCache
- OS 调优三件套:
vm.swappiness=1、文件句柄 10w+、XFS noatime - 运维三大命令:reassign-partitions(带 throttle!)、leader-election、consumer-groups reset
- 故障演练:每季度一次,让团队对故障有肌肉记忆
下一章我们进入「踩坑案例集」——12 个真实事故复盘,每一个都可能让你少加 N 个班。
🎬 可视化演示
演示加载缓慢或样式异常?点此在新标签页打开 ↗
💻 示例代码
yaml
# JMX Exporter 规则示例(Kafka 3.x)
# 用法:
# export KAFKA_OPTS="-javaagent:/opt/jmx_exporter/jmx_prometheus_javaagent.jar=7071:$(pwd)/jmx_exporter_config.yml"
# bin/kafka-server-start.sh config/server.properties
# 然后访问 http://broker:7071/metrics 即可看到 Prometheus 文本格式
lowercaseOutputName: true
lowercaseOutputLabelNames: true
# 不要把所有 MBean 都暴露,只取关键的,体积才合理
whitelistObjectNames:
- "kafka.server:type=BrokerTopicMetrics,*"
- "kafka.server:type=ReplicaManager,*"
- "kafka.server:type=ReplicaFetcherManager,*"
- "kafka.server:type=KafkaRequestHandlerPool,*"
- "kafka.network:type=SocketServer,*"
- "kafka.network:type=RequestMetrics,*"
- "kafka.controller:type=KafkaController,*"
- "kafka.controller:type=ControllerStats,*"
- "kafka.log:type=LogManager,*"
- "kafka.log:type=LogCleanerManager,*"
- "kafka.cluster:type=Partition,*"
- "java.lang:type=GarbageCollector,name=*"
- "java.lang:type=Memory"
- "java.lang:type=OperatingSystem"
rules:
# ========== Broker 健康 ==========
- pattern: 'kafka.server<type=ReplicaManager, name=(.+)><>Value'
name: kafka_server_replicamanager_$1
type: GAUGE
- pattern: 'kafka.controller<type=KafkaController, name=(.+)><>Value'
name: kafka_controller_kafkacontroller_$1
type: GAUGE
- pattern: 'kafka.controller<type=ControllerStats, name=(.+)><>(Count|MeanRate|OneMinuteRate|FiveMinuteRate)'
name: kafka_controller_controllerstats_$1_$2
type: GAUGE
# ========== 流量 ==========
- pattern: 'kafka.server<type=BrokerTopicMetrics, name=(.+), topic=(.+)><>(Count|OneMinuteRate)'
name: kafka_server_brokertopicmetrics_$1_$3
type: GAUGE
labels:
topic: "$2"
- pattern: 'kafka.server<type=BrokerTopicMetrics, name=(.+)><>(Count|OneMinuteRate)'
name: kafka_server_brokertopicmetrics_$1_$2
type: GAUGE
# ========== 请求队列 / 处理 ==========
- pattern: 'kafka.network<type=RequestMetrics, name=(.+), request=(.+)><>(Count|Mean|99thPercentile|999thPercentile)'
name: kafka_network_requestmetrics_$1_$3
type: GAUGE
labels:
request: "$2"
- pattern: 'kafka.server<type=KafkaRequestHandlerPool, name=RequestHandlerAvgIdlePercent><>OneMinuteRate'
name: kafka_server_kafkarequesthandlerpool_requesthandleravgidlepercent
type: GAUGE
- pattern: 'kafka.network<type=SocketServer, name=NetworkProcessorAvgIdlePercent><>Value'
name: kafka_network_socketserver_networkprocessoravgidlepercent
type: GAUGE
# ========== 副本 / Fetcher ==========
- pattern: 'kafka.server<type=ReplicaFetcherManager, name=MaxLag, clientId=(.+)><>Value'
name: kafka_server_replicafetchermanager_maxlag
type: GAUGE
labels:
client: "$1"
- pattern: 'kafka.cluster<type=Partition, name=UnderReplicated, topic=(.+), partition=(.+)><>Value'
name: kafka_cluster_partition_underreplicated
type: GAUGE
labels:
topic: "$1"
partition: "$2"
# ========== 日志 / Compaction ==========
- pattern: 'kafka.log<type=LogManager, name=LogDirectoryOffline, logDirectory=(.+)><>Value'
name: kafka_log_logmanager_offline_logdir
type: GAUGE
labels:
logdir: "$1"
- pattern: 'kafka.log<type=Log, name=Size, topic=(.+), partition=(.+)><>Value'
name: kafka_log_log_size
type: GAUGE
labels:
topic: "$1"
partition: "$2"
# ========== JVM ==========
- pattern: 'java.lang<type=Memory><HeapMemoryUsage>used'
name: jvm_memory_heap_used_bytes
type: GAUGE
- pattern: 'java.lang<type=GarbageCollector, name=(.+)><>CollectionCount'
name: jvm_gc_collection_count
type: COUNTER
labels:
gc: "$1"
- pattern: 'java.lang<type=GarbageCollector, name=(.+)><>CollectionTime'
name: jvm_gc_collection_time_ms
type: COUNTER
labels:
gc: "$1"
- pattern: 'java.lang<type=OperatingSystem><>(SystemCpuLoad|ProcessCpuLoad|FreePhysicalMemorySize)'
name: jvm_os_$1
type: GAUGEpython
#!/usr/bin/env python3
"""
Consumer Lag 监控 + Webhook 告警。
依赖:
pip install confluent-kafka requests
示例:
BOOTSTRAP=localhost:9092 LAG_THRESHOLD=10000 \
WEBHOOK_URL=https://hooks.example.com/xxx python lag_alert.py
"""
import os
import time
import requests
from confluent_kafka.admin import AdminClient, ConsumerGroupTopicPartitions
from confluent_kafka import Consumer, TopicPartition
BOOTSTRAP = os.getenv("BOOTSTRAP", "localhost:9092")
LAG_THRESHOLD = int(os.getenv("LAG_THRESHOLD", "10000"))
INTERVAL = int(os.getenv("INTERVAL", "30"))
WEBHOOK = os.getenv("WEBHOOK_URL", "")
def list_groups(admin):
fut = admin.list_consumer_groups(request_timeout=10)
res = fut.result()
return [g.group_id for g in res.valid if not g.group_id.startswith("_")]
def get_group_lag(group_id):
admin = AdminClient({"bootstrap.servers": BOOTSTRAP})
consumer = Consumer({
"bootstrap.servers": BOOTSTRAP,
"group.id": group_id,
"enable.auto.commit": False,
})
try:
fut = admin.list_consumer_group_offsets(
[ConsumerGroupTopicPartitions(group_id, None)]
)
offsets_res = list(fut.values())[0].result()
committed = offsets_res.topic_partitions
rows = []
for tp in committed:
if tp.error:
continue
tp_q = TopicPartition(tp.topic, tp.partition)
low, high = consumer.get_watermark_offsets(tp_q, timeout=5, cached=False)
cur = tp.offset if tp.offset >= 0 else low
lag = max(0, high - cur)
rows.append((tp.topic, tp.partition, cur, high, lag))
return rows
finally:
consumer.close()
def alert(group, summary, details):
text = "[Kafka Lag Alert] group={}\n{}\n\n{}".format(group, summary, details)
print(text)
if WEBHOOK:
try:
requests.post(
WEBHOOK,
json={"msgtype": "text", "text": {"content": text}},
timeout=5,
)
except Exception as e:
print(" webhook fail:", e)
def main():
print("Lag Monitor: bootstrap={} threshold={} interval={}s".format(
BOOTSTRAP, LAG_THRESHOLD, INTERVAL))
admin = AdminClient({"bootstrap.servers": BOOTSTRAP})
while True:
try:
groups = list_groups(admin)
print("\n[{}] checking {} groups".format(
time.strftime("%H:%M:%S"), len(groups)))
for g in groups:
rows = get_group_lag(g)
if not rows:
continue
total_lag = sum(r[4] for r in rows)
worst = max(rows, key=lambda r: r[4])
print(" {:30s} total={:>10d} worst={}-{} lag={}".format(
g, total_lag, worst[0], worst[1], worst[4]))
if worst[4] >= LAG_THRESHOLD:
top = sorted(rows, key=lambda r: -r[4])[:5]
detail = "\n".join(
" {}-{}: cur={} end={} lag={}".format(t, p, c, e, l)
for (t, p, c, e, l) in top
)
alert(
group=g,
summary="max-lag={} total={}".format(worst[4], total_lag),
details="Top 5 lagged partitions:\n" + detail,
)
except Exception as e:
print(" ERR:", e)
time.sleep(INTERVAL)
if __name__ == "__main__":
try:
main()
except KeyboardInterrupt:
print("\nbye")yaml
# Prometheus 抓取配置片段(合并到 prometheus.yml 的 scrape_configs:)
# 启动:docker run --network host -v $(pwd)/prometheus.yml:/etc/prometheus/prometheus.yml prom/prometheus
global:
scrape_interval: 15s
evaluation_interval: 15s
external_labels:
cluster: prod-1
region: cn-bj-1
rule_files:
- "kafka_alerts.yml"
alerting:
alertmanagers:
- static_configs:
- targets: ['alertmanager:9093']
scrape_configs:
# ===== Broker JMX Exporter =====
- job_name: kafka-broker
metrics_path: /metrics
static_configs:
- targets:
- kafka1:7071
- kafka2:7071
- kafka3:7071
labels:
role: broker
relabel_configs:
- source_labels: [__address__]
target_label: instance
regex: '([^:]+):.*'
replacement: '$1'
# ===== kafka-exporter (Lag 专用) =====
# https://github.com/danielqsj/kafka_exporter
- job_name: kafka-exporter
static_configs:
- targets: ['kafka-exporter:9308']
# ===== Schema Registry =====
- job_name: schema-registry
metrics_path: /metrics
static_configs:
- targets: ['schema-registry:5556']
# ===== Connect Worker =====
- job_name: kafka-connect
metrics_path: /metrics
static_configs:
- targets: ['connect:7072']
# =====================================================
# 一份开箱即用的告警规则(可单独保存为 kafka_alerts.yml)
# =====================================================
# groups:
# - name: kafka.rules
# rules:
# - alert: KafkaUnderReplicated
# expr: sum(kafka_server_replicamanager_underreplicatedpartitions) by (instance) > 0
# for: 5m
# labels: { severity: warning, team: data-infra }
# annotations:
# summary: "{{ $labels.instance }} 副本不全 ({{ $value }} 分区)"
#
# - alert: KafkaOfflinePartitions
# expr: sum(kafka_controller_kafkacontroller_offlinepartitionscount) > 0
# for: 1m
# labels: { severity: critical, page: oncall }
# annotations:
# summary: "Kafka 集群有 Offline 分区!"
#
# - alert: KafkaActiveControllerWrong
# expr: sum(kafka_controller_kafkacontroller_activecontrollercount) != 1
# for: 1m
# labels: { severity: critical }
#
# - alert: KafkaUncleanLeaderElection
# expr: sum(rate(kafka_controller_controllerstats_uncleanleaderelectionspersec_count[5m])) > 0
# labels: { severity: critical }
# annotations: { summary: "出现 Unclean Leader Election,可能丢数据!" }
#
# - alert: KafkaConsumerLagHigh
# expr: kafka_consumergroup_lag > 10000
# for: 10m
# labels: { severity: warning }
#
# - alert: KafkaRequestHandlerLow
# expr: kafka_server_kafkarequesthandlerpool_requesthandleravgidlepercent < 0.2
# for: 10m
# labels: { severity: warning }
# annotations: { summary: "{{ $labels.instance }} I/O 线程繁忙,考虑加 num.io.threads" }bash
#!/usr/bin/env bash
# 分区重分配封装脚本:生成 → 执行 → 验证 → 取消限速
# 用法:
# ./reassign_helper.sh generate "topic1,topic2" "1,2,3,4" # 生成方案
# ./reassign_helper.sh execute reassign.json 50000000 # 执行(限速 50MB/s)
# ./reassign_helper.sh verify reassign.json # 验证
# ./reassign_helper.sh cancel-throttle reassign.json # 单独取消限速
set -euo pipefail
BS=${BS:-localhost:9092}
KAFKA_HOME=${KAFKA_HOME:-/opt/kafka}
SCRIPT="$KAFKA_HOME/bin/kafka-reassign-partitions.sh"
usage() {
cat <<EOF
Usage:
$0 generate <topic_csv> <broker_csv>
生成 reassign.json(基于 broker-list 平均分布)
$0 execute <reassign.json> [throttle_bytes_per_sec]
执行(默认限速 50MB/s)
$0 verify <reassign.json>
验证当前进度
$0 cancel-throttle <reassign.json>
手动取消限速(verify 完成后会自动取消,但失败时可手动)
EOF
exit 1
}
cmd=${1:-}
case "$cmd" in
generate)
topics=$2; brokers=$3
# 生成 topics-to-move JSON
TJSON=$(mktemp)
echo -n '{"topics":[' > "$TJSON"
IFS=',' read -ra arr <<< "$topics"
for i in "${!arr[@]}"; do
[[ $i -gt 0 ]] && echo -n ',' >> "$TJSON"
echo -n "{\"topic\":\"${arr[$i]}\"}" >> "$TJSON"
done
echo '] ,"version":1}' >> "$TJSON"
echo "==> 生成方案 (broker-list=$brokers)"
"$SCRIPT" --bootstrap-server "$BS" \
--topics-to-move-json-file "$TJSON" \
--broker-list "$brokers" \
--generate
echo ""
echo "✅ 把上面 'Proposed partition reassignment configuration' 段落保存为 reassign.json"
;;
execute)
plan=$2; throttle=${3:-50000000}
echo "==> 执行重分配,限速 = $throttle bytes/s ($((throttle/1024/1024)) MB/s)"
"$SCRIPT" --bootstrap-server "$BS" \
--reassignment-json-file "$plan" \
--execute --throttle "$throttle"
echo "✅ 已提交,下一步用 verify 查看进度"
;;
verify)
plan=$2
echo "==> 验证进度"
"$SCRIPT" --bootstrap-server "$BS" \
--reassignment-json-file "$plan" \
--verify
;;
cancel-throttle)
plan=$2
echo "==> 取消限速"
"$SCRIPT" --bootstrap-server "$BS" \
--reassignment-json-file "$plan" \
--verify
echo "(verify 会自动清理;如仍残留请用 kafka-configs.sh 手动 delete)"
;;
*) usage ;;
esacjmx_exporter_config.yml ↗ · lag_alert.py ↗ · prometheus_scrape.yml ↗ · reassign_helper.sh ↗