Skip to content

第 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 的分区数= 0Follower 掉队 / 网络延迟 / 磁盘慢;查 replica.lag.time.max.ms 与 ReplicaFetcher Lag
OfflinePartitionsCount没有 Leader 的分区数(最致命= 0Broker 全部宕机或 Controller 异常;立刻看 controller.log
ActiveControllerCount当前 Controller 数整集群恰为 1=0:没人当 Controller,集群脑死;>1:脑裂,立刻人工介入
LeaderElectionRateAndTimeMsLeader 选举速率与耗时平时 ≈ 0突增 = Broker 抖动 / 网络问题,看 state-change.log
UncleanLeaderElectionsPerSec不干净选举速率(会丢数据= 0已发生丢数据;检查是否误开 unclean.leader.election.enable=true

1.2 流量与性能类

指标含义健康阈值异常排查方向
MessagesInPerSec每秒进入消息数(每 Topic / 全局)看业务基线突降 = 上游异常;突增 = 上游打满
BytesInPerSec每秒入流量< 网卡带宽 70%接近上限要扩容
BytesOutPerSec每秒出流量< 网卡带宽 70%同上;通常 OutPerSec 是 InPerSec × 副本数 + Consumer fanout
RequestHandlerAvgIdlePercentI/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 类

指标含义健康阈值异常排查方向
ISRShrinksPerSecISR 收缩速率≈ 0Follower 掉队 / 网络抖动
ISRExpandsPerSecISR 扩张速率≈ 0与 Shrink 配对出现是「ISR 抖动」
ReplicaFetcher MaxLag副本同步最大滞后几千以内持续高 = 该 Follower 跟不上,检查带宽 / GC
LeaderCountLeader 数各 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平均请求延迟< 50msp99 单独看
request-latency-p99p99 延迟< 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.properties

kafka.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名称说明
7589Kafka Exporter OverviewTopic / Group / Lag 总览
11962Kafka (JMX Exporter)Broker 全维度
21078Kafka Cluster OverviewConfluent 风格
18276Strimzi KafkaK8s 部署专用

📌 直接在 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 实时盯关键词:

bash
tail -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.threads
  • NetworkProcessorAvgIdlePercent 长期 < 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 != 15 分钟内有人响应
P1 严重UnderReplicated > 0 持续 5 分钟 / Lag > 阈值30 分钟内响应
P2 警告RequestHandlerIdle < 20% / GC 时间高1 小时内查看
P3 通知Topic 增长率异常 / 磁盘 60%上班时间处理

7.2 三个「上线必须验证」的指标

  1. UnderReplicatedPartitions = 0:副本健康
  2. ActiveControllerCount = 1(全集群):Controller 健康
  3. UncleanLeaderElectionsPerSec = 0:无数据丢失

7.3 一定要看的「业务指标」

  • 各 Topic 的 实际吞吐(消息数 / 字节数)
  • 各 Consumer Group 的 Lag 趋势
  • 各 Topic 的 平均消息大小(突增意味着业务异常或日志爆炸)
  • 生产者错误率(按错误类型分类)

8. 故障演练(运维成熟度的标尺)

每季度做一次:

演练操作期望结果
Broker 宕机kill -9 一个 BrokerLeader 在 30 秒内切走,业务无感
磁盘满dd 把 log 目录填满该 Broker 拒收写入但不 panic,迁出后恢复
网络分区iptables drop 与某 Broker 的连接ISR 收缩,业务降级但不崩溃
Controller 宕机kill 当前 Controller 节点选举新 Controller,期间元数据变更阻塞 < 30 秒
Schema Registry 挂stop SRProducer/Consumer 走本地缓存,新 Schema 注册失败但已有业务正常

9. 横向对比:与其他 MQ 的运维差异

维度KafkaRabbitMQRocketMQPulsar
元数据KRaft / ZKMnesia / 内存NameServerZK + BookKeeper
主要指标UnderReplicated / Lagqueue length / consumersbroker offset / consume lagbookie ledger
重分配reassign-partitions(手动)队列绑定改路由调整队列分布bundle 自动 reassign
重启代价较高(PageCache 重建)中(持久化队列恢复慢)低(Broker 无状态)
单机扩展加 broker + reassign加 node + mirror加 broker + topic 扩展加 broker / 加 bookie 解耦

10. 本章面试高频题

Q1:判断 Kafka 集群健康,最先看哪几个指标?

答案(背 5 个)

  1. OfflinePartitionsCount —— 必须 = 0(否则有数据完全无法读写)
  2. UnderReplicatedPartitions —— 应 = 0(持续大于 0 说明副本掉队)
  3. ActiveControllerCount —— 整集群恰为 1(>1 = 脑裂,=0 = 集群瘫痪)
  4. UncleanLeaderElectionsPerSec —— = 0(>0 说明发生了「丢数据」式选举)
  5. RequestHandlerAvgIdlePercent —— > 30%(< 20% 说明 I/O 线程不够)

加分:能解释每个指标的物理含义、阈值由来、常见根因。


Q2:Consumer Lag 一直在涨,怎么排查?

结构化答

  1. 先确认 Lag 真涨kafka-consumer-groups.sh --describe 看具体分区的 Lag
  2. 定位是上游写多 / 下游消费慢:对比 MessagesInPerSec 与 Consumer records-consumed-rate
  3. 如是消费慢
    • 业务处理慢 → 看应用日志、数据库慢查询、外部依赖
    • 单线程消费 → 用 max.poll.records + 多线程处理
    • Rebalance 频繁 → 调 max.poll.interval.ms / session.timeout.ms
  4. 结构性原因:分区数 < 消费者期望并发数 → 加分区
  5. 临时止血:扩容消费者实例(不超过分区数)

加分项:提到 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 时怎么处理?

答案分阶段

  1. controller.log / server.log:定位是哪个 Broker / Topic / Partition
  2. 看 ReplicaFetcher Lag:是不是某个 Follower 跟不上
  3. 常见原因
    • Broker 宕机 / 网络隔离 → 重启或修网络
    • Broker 磁盘满 → 清日志 / 加盘 / 迁移分区
    • GC 暂停长 → 调 JVM
    • 副本同步限速 (replica.fetch.max.bytes / num.replica.fetchers) 配置过低 → 调高
  4. 临时缓解:把该 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: GAUGE
python
#!/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 ;;
esac

jmx_exporter_config.yml ↗ · lag_alert.py ↗ · prometheus_scrape.yml ↗ · reassign_helper.sh ↗