Kafka 运维监控与故障恢复:JMX 指标、Lag 监控与分区重分配

Kafka 生产集群运维:JMX 核心指标、Consumer Lag 监控告警、分区重分配、Broker 扩缩容与数据恢复

Kafka 生产集群的稳定运行离不开完善的监控体系与规范的运维流程。一个未经监控的 Kafka 集群如同在黑暗中驾驶,分区复制延迟、消费者堆积、Broker 状态异常等问题随时可能引发服务故障。本篇文章系统地梳理了 Kafka 运维的核心维度,涵盖 JMX 关键指标解读、Consumer Lag 实时监控、分区重分配操作、Broker 扩缩容、数据恢复策略以及 Prometheus + Grafana 监控方案,帮助运维人员建立可观测、可预警、可恢复的生产级 Kafka 运维体系。


一、JMX 核心指标

Kafka 暴露了丰富的 JMX(Java Management Extensions)指标,运维人员可以通过这些指标洞察集群的健康状况。理解核心指标的含义和告警阈值是高效运维的基础。

1.1 关键服务端指标

UnderReplicatedPartitions

该指标统计当前存在副本同步延迟的分区数量。正常情况下,每个分区的所有副本都应该处于同步状态(在 ISR 列表中)。当 Broker 故障、网络抖动或副本落后领导者过多时,Follower 副本会被踢出 ISR,导致该指标上升。

# 通过 JMX 查看 UnderReplicatedPartitions
jmxterm << JMXEOF
open localhost:9999
get -b kafka.server:type=ReplicaManager,name=UnderReplicatedPartitions -a Value
JMXEOF

告警阈值

  • 值大于 0 持续超过 5 分钟:触发严重告警
  • 值大于 0 但快速恢复(1-2 分钟内归零):记录为事件但不告警

OfflinePartitions

表示当前有领导者副本不可用的分区数量。当分区的 Leader 副本所在 Broker 宕机,且没有可用的 ISR 副本可接管时,该分区就会进入 Offline 状态,导致生产者和消费者都无法读写该分区。

# 查看 OfflinePartitions
jmxterm << JMXEOF
open localhost:9999
get -b kafka.server:type=ReplicaManager,name=OfflinePartitions -a Value
JMXEOF

告警阈值

  • 任何大于 0 的情况都应立即触发严重告警,因为这代表分区不可用
  • 该指标持续非零超过 30 秒即为重大故障

ISRShrink / ISRExpand

  • ISRShrink:统计 ISR 列表缩小的次数,通常由于 Follower 副本落后过多被移除 ISR
  • ISRExpand:统计 ISR 列表扩展的次数,通常由于 Follower 副本追赶上 Leader
# 查看 ISRShrink 指标
jmxterm << JMXEOF
open localhost:9999
get -b kafka.server:type=ReplicaManager,name=ISRShrinksPerSec -a Count
JMXEOF

告警阈值

  • ISRShrink 速率持续大于 0 且 ISRExpand 速率接近 0:严重告警
  • 频繁出现 ISRShrink 说明集群存在性能问题或网络不稳定

ActiveControllerCount

统计当前活跃的 Controller 数量。Kafka 集群中应该只有一个 Controller。当 ActiveControllerCount 不等于 1 时,说明 Controller 选举出现异常。

# 查看 ActiveControllerCount
jmxterm << JMXEOF
open localhost:9999
get -b kafka.controller:type=KafkaController,name=ActiveControllerCount -a Value
JMXEOF

告警阈值

  • 值不等于 1:立即触发严重告警
  • 等于 0 表示集群没有 Controller,无法执行管理操作
  • 大于 1 可能出现 “脑裂” 情况,极为罕见但非常严重

UncleanLeaderElectionsPerSec

统计每分钟发生的非干净 Leader 选举次数。非干净选举指的是从不在 ISR 中的副本中选举 Leader,这会导致数据丢失(因为这些副本可能缺少部分消息)。

# 查看 UncleanLeaderElectionsPerSec
jmxterm << JMXEOF
open localhost:9999
get -b kafka.controller:type=ControllerStats,name=UncleanLeaderElectionsPerSec -a Count
JMXEOF

告警阈值

  • 值大于 0:立即触发严重告警
  • 非干净选举默认被禁用(unclean.leader.election.enable=false),任何非零值都说明配置被覆盖或出现了极端故障

1.2 指标总览表

指标名称MBean 路径正常值告警阈值
UnderReplicatedPartitionskafka.server:type=ReplicaManager,name=UnderReplicatedPartitions0> 0 持续 5min
OfflinePartitionskafka.server:type=ReplicaManager,name=OfflinePartitions0> 0
ActiveControllerCountkafka.controller:type=KafkaController,name=ActiveControllerCount1!= 1
UncleanLeaderElectionskafka.controller:type=ControllerStats,name=UncleanLeaderElectionsPerSec0> 0
ISRShrinksPerSeckafka.server:type=ReplicaManager,name=ISRShrinksPerSec0> 0 持续 3min

二、Broker 健康检查

Broker 是 Kafka 集群的基本节点,每个 Broker 的状态直接影响整个集群的可用性。系统性的健康检查应覆盖网络、磁盘、内存、CPU 和 Kafka 进程本身。

2.1 进程与端口检查

# 检查 Kafka 进程是否存活
systemctl status kafka
pgrep -f "kafka.Kafka"

# 检查监听端口(默认 9092 为数据端口,9093 为内部通信端口,9999 为 JMX 端口)
ss -tlnp | grep -E '9092|9093|9999'

# 检查 Broker 是否能响应元数据请求
kafka-broker-api-versions.sh --bootstrap-server localhost:9092

2.2 磁盘与水珠线检查

Kafka 严重依赖磁盘 IO,磁盘空间不足是导致集群故障的常见原因。

# 检查数据目录磁盘使用率
df -h /var/lib/kafka-logs

# 检查日志段保留策略执行情况
ls -lh /var/lib/kafka-logs/*/ | head -20

# Kafka 自带的日志段大小检查脚本
kafka-log-dirs.sh --bootstrap-server localhost:9092 \
  --describe \
  --topic-list my-topic

磁盘告警阈值

  • 使用率超过 70%:预警,准备清理或扩容
  • 使用率超过 85%:严重告警,立即处理
  • 使用率超过 90%:紧急告警,Kafka 写入可能失败

2.3 网络检查

# 检查 Broker 间网络延迟
for broker in broker1:9092 broker2:9092 broker3:9092; do
  echo "Testing $broker..."
  kafka-broker-api-versions.sh --bootstrap-server $broker
done

# 使用 telnet 检查连通性
echo "exit" | telnet broker2 9092 2>/dev/null | grep Connected

2.4 Zookeeper 连接检查

# 检查 Kafka 在 Zookeeper 中的注册状态
zookeeper-shell.sh zk1:2181 <<< "ls /brokers/ids"

# 检查每个 Broker 的元数据
zookeeper-shell.sh zk1:2181 <<< "get /brokers/ids/1"

正常输出应显示所有 Broker ID 列表。如果某个 Broker ID 缺失,说明该 Broker 未成功注册到 Zookeeper。


三、Consumer Lag 监控

Consumer Lag 表示消费者当前消费位置与分区最新消息位置的差距,是 Kafka 运维中最重要的监控指标之一。Lag 持续增大意味着消费速度跟不上生产速度,可能导致消息处理延迟甚至内存溢出。

3.1 使用 kafka-consumer-groups.sh 查看 Lag

kafka-consumer-groups.sh 是 Kafka 自带的消费者组管理工具,可以查看所有消费者组的消费状态。

# 列出所有消费者组
kafka-consumer-groups.sh \
  --bootstrap-server localhost:9092 \
  --list

# 查看指定消费者组的详细消费状态
kafka-consumer-groups.sh \
  --bootstrap-server localhost:9092 \
  --describe \
  --group order-service-group

典型输出:

GROUP                  TOPIC           PARTITION  CURRENT-OFFSET  LOG-END-OFFSET  LAG   CONSUMER-ID
order-service-group    orders          0          1284500         1284502         2     consumer-1
order-service-group    orders          1          983400          983400          0     consumer-2
order-service-group    orders          2          1567800         1567890         90    consumer-3

字段说明

  • CURRENT-OFFSET:消费者当前提交的消费位置
  • LOG-END-OFFSET:分区中最新消息的偏移量
  • LAG:未消费的消息数量
  • CONSUMER-ID:消费该分区的消费者实例 ID

3.2 批量 Lag 监控脚本

#!/bin/bash
# monitor-lag.sh - 批量监控所有消费者组的 Lag

BOOTSTRAP_SERVER="localhost:9092"
ALERT_THRESHOLD=1000
SEVERE_THRESHOLD=10000

echo "=== Consumer Lag Report $(date) ==="
echo ""

kafka-consumer-groups.sh \
  --bootstrap-server $BOOTSTRAP_SERVER \
  --list | while read group; do

  echo "Group: $group"

  kafka-consumer-groups.sh \
    --bootstrap-server $BOOTSTRAP_SERVER \
    --describe \
    --group "$group" 2>/dev/null | \
    awk -v alert=$ALERT_THRESHOLD -v severe=$SEVERE_THRESHOLD '
    NR>1 {
      lag = $5
      partition = $3
      if (lag >= severe) {
        print "  [SEVERE] Partition " partition " LAG=" lag
      } else if (lag >= alert) {
        print "  [ALERT]  Partition " partition " LAG=" lag
      }
    }'
done

3.3 使用 Burrow 进行高级 Lag 监控

Burrow 是 LinkedIn 开源的 Kafka Consumer Lag 监控工具,它可以评估消费状态是否健康,而不仅仅是显示当前 Lag 数值。

Burrow 安装与配置

# 下载 Burrow
git clone https://github.com/linkedin/Burrow.git
cd Burrow
go build

# 创建配置文件
mkdir -p /etc/burrow
cat > /etc/burrow/burrow.toml << 'CONF'
[general]
logdir=/var/log/burrow
logconfig=/etc/burrow/logconfig.xml

[zookeeper]
servers=["zk1:2181", "zk2:2181", "zk3:2181"]
timeout=6

[kafka.kafka-cluster]
servers=["broker1:9092", "broker2:9092", "broker3:9092"]

[consumer.kafka-cluster]
servers=["broker1:9092", "broker2:9092", "broker3:9092"]

[httpserver.default]
address=":8000"
CONF

Burrow 状态评估

Burrow 将消费者状态分为四种:

状态含义处理建议
OK消费正常,Lag 保持稳定或递减无需处理
WARNING消费略慢,Lag 缓慢增长关注趋势
ERR消费明显滞后,Lag 快速增长立即排查
STOP消费完全停止紧急处理
# 查询某个消费者组的状态
curl -s http://localhost:8000/v3/kafka/kafka-cluster/consumer/my-group/lag | jq .

Burrow 告警集成

# 提取处于 ERR 或 STOP 状态的消费者组
curl -s http://localhost:8000/v3/kafka/kafka-cluster/consumer | jq -r '
  .consumers[] |
  select(.status == "ERR" or .status == "STOP") |
  "ALERT: Group \(.name) is in \(.status) state"
'

四、告警阈值设置

合理的告警阈值应基于业务特点和历史数据设定,避免 “狼来了” 式的无效告警。

4.1 分层告警策略

# alert-rules.yml - Prometheus 告警规则示例
groups:
  - name: kafka-alerts
    rules:
      # 严重级别:分区不可用
      - alert: KafkaOfflinePartitions
        expr: kafka_server_replicamanager_offlinepartitioncount > 0
        for: 30s
        labels:
          severity: critical
        annotations:
          summary: "Kafka has offline partitions"
          description: "{{ $value }} partitions are offline"

      # 严重级别:副本同步异常
      - alert: KafkaUnderReplicatedPartitions
        expr: kafka_server_replicamanager_underreplicatedpartitions > 0
        for: 5m
        labels:
          severity: critical
        annotations:
          summary: "Kafka under replicated partitions detected"
          description: "{{ $value }} partitions are under-replicated"

      # 警告级别:Controller 异常
      - alert: KafkaActiveController
        expr: kafka_controller_kafkacontroller_activecontrollercount != 1
        for: 1m
        labels:
          severity: warning
        annotations:
          summary: "Kafka controller anomaly"
          description: "Active controller count is {{ $value }}"

      # 警告级别:Consumer Lag 过大
      - alert: KafkaConsumerLagHigh
        expr: kafka_consumergroup_lag > 10000
        for: 5m
        labels:
          severity: warning
        annotations:
          summary: "High consumer lag"
          description: "Consumer group lag is {{ $value }}"

      # 信息级别:磁盘使用率
      - alert: KafkaDiskSpaceLow
        expr: kafka_log_log_size / kafka_log_log_capacity > 0.7
        for: 10m
        labels:
          severity: info
        annotations:
          summary: "Kafka disk space over 70%"
          description: "Disk usage is {{ $value | humanizePercentage }}"

4.2 告警通知配置示例

# alertmanager.yml 配置
route:
  group_by: ['alertname', 'severity']
  group_wait: 30s
  group_interval: 5m
  repeat_interval: 4h
  receiver: 'kafka-team'

receivers:
  - name: 'kafka-team'
    slack_configs:
      - channel: '#kafka-alerts'
        title: '{{ range .Alerts }}{{ .Annotations.summary }}{{ end }}'
        text: '{{ range .Alerts }}{{ .Annotations.description }}{{ end }}'
    pagerduty_configs:
      - service_key: '<your-pagerduty-key>'
        severity: '{{ .GroupLabels.severity }}'

五、分区重分配

当集群扩容、Broker 故障恢复或原有分区分布不均时,需要对分区进行重分配(Reassignment)。分区重分配会涉及大量数据迁移,应谨慎规划执行时机。

5.1 生成重分配计划

# 查看当前分区分布
kafka-reassign-partitions.sh \
  --bootstrap-server localhost:9092 \
  --describe \
  --topics-to-move-json-file topics.json

首先创建 topics.json,声明需要重分配的 Topic:

{
  "topics": [
    {"topic": "orders"},
    {"topic": "payments"},
    {"topic": "notifications"}
  ],
  "version": 1
}

然后使用 --generate 自动生成重分配方案:

kafka-reassign-partitions.sh \
  --bootstrap-server localhost:9092 \
  --topics-to-move-json-file topics.json \
  --broker-list "1,2,3,4,5" \
  --generate

输出包含两部分:

  • Current partition replica assignment:当前分配状态
  • Proposed partition reassignment configuration:建议的重分配方案

5.2 执行重分配

将建议方案保存为 reassignment.json,然后执行:

kafka-reassign-partitions.sh \
  --bootstrap-server localhost:9092 \
  --reassignment-json-file reassignment.json \
  --execute

5.3 监控重分配进度

# 验证重分配进度(重复执行查看状态)
kafka-reassign-partitions.sh \
  --bootstrap-server localhost:9092 \
  --reassignment-json-file reassignment.json \
  --verify

输出示例:

Status of partition reassignment:
Reassignment of partition orders-0 is still in progress
Reassignment of partition orders-1 completed successfully
Reassignment of partition orders-2 is still in progress

5.4 限流重分配

重分配期间大量数据复制可能影响正常业务流量。Kafka 支持通过 throttle 参数限制重分配速率:

# 执行时设置限流(例如限制为 100MB/s)
kafka-reassign-partitions.sh \
  --bootstrap-server localhost:9092 \
  --reassignment-json-file reassignment.json \
  --execute \
  --throttle 104857600

# 重分配完成后取消限流
kafka-reassign-partitions.sh \
  --bootstrap-server localhost:9092 \
  --reassignment-json-file reassignment.json \
  --verify

5.5 优先副本选举

重分配完成后,建议执行优先副本选举,确保 Leader 均匀分布:

kafka-leader-election.sh \
  --bootstrap-server localhost:9092 \
  --election-type preferred \
  --all-topic-partitions

六、Broker 扩缩容

6.1 Broker 扩容

扩容 Broker 是提升集群存储和吞吐能力的常见操作。

步骤一:准备新 Broker 节点

# 复制 Kafka 配置到新节点
scp /opt/kafka/config/server.properties new-broker:/opt/kafka/config/

# 修改新 Broker 配置(确保 broker.id 唯一)
cat >> /opt/kafka/config/server.properties << 'PROPS'
broker.id=4
listeners=PLAINTEXT://:9092
log.dirs=/var/lib/kafka-logs
zookeeper.connect=zk1:2181,zk2:2181,zk3:2181
PROPS

步骤二:启动新 Broker

systemctl start kafka
# 或
/opt/kafka/bin/kafka-server-start.sh -daemon /opt/kafka/config/server.properties

步骤三:验证新 Broker 加入

# 确认 Broker 已注册到 Zookeeper
zookeeper-shell.sh zk1:2181 <<< "ls /brokers/ids"
# 应输出包含新 broker.id,例如: [1, 2, 3, 4]

# 验证新 Broker 可提供服务
kafka-broker-api-versions.sh --bootstrap-server new-broker:9092

步骤四:迁移数据到新 Broker

新 Broker 加入后,默认不会自动接收已有 Topic 的数据。需要执行分区重分配,将部分分区迁移到新 Broker:

# 生成包含新 Broker 的重分配方案
kafka-reassign-partitions.sh \
  --bootstrap-server localhost:9092 \
  --topics-to-move-json-file topics.json \
  --broker-list "1,2,3,4" \
  --generate > plan.txt

# 提取建议方案并执行
grep -A 100 'Proposed partition reassignment' plan.txt | tail -n +2 > reassignment.json
kafka-reassign-partitions.sh \
  --bootstrap-server localhost:9092 \
  --reassignment-json-file reassignment.json \
  --execute

6.2 Broker 缩容

缩容操作需要先将目标 Broker 上的数据迁移到其他 Broker,然后安全下线。

# 步骤一:生成分区重分配方案(排除即将下线的 Broker)
kafka-reassign-partitions.sh \
  --bootstrap-server localhost:9092 \
  --topics-to-move-json-file topics.json \
  --broker-list "1,2,3" \
  --generate > plan.txt

# 步骤二:执行重分配(将 Broker 4 的数据迁出)
grep -A 100 'Proposed' plan.txt | tail -n +2 > reassignment.json
kafka-reassign-partitions.sh \
  --bootstrap-server localhost:9092 \
  --reassignment-json-file reassignment.json \
  --execute

# 步骤三:监控直到所有分区迁移完成
kafka-reassign-partitions.sh \
  --bootstrap-server localhost:9092 \
  --reassignment-json-file reassignment.json \
  --verify

# 步骤四:安全关闭目标 Broker
kafka-server-stop.sh
# 或
systemctl stop kafka

# 步骤五:从 Zookeeper 确认 Broker 已移除
zookeeper-shell.sh zk1:2181 <<< "ls /brokers/ids"

七、数据恢复

7.1 日志恢复

Kafka 的日志文件是顺序写入的,在 Broker 非正常关闭(如断电、OOM 被杀)后,可能需要修复日志文件。

# 使用 Kafka 自带的日志恢复工具
/opt/kafka/bin/kafka-recovery-point-offset-checkpoint.sh \
  /var/lib/kafka-logs

# 检查特定 Topic 分区的日志完整性
kafka-dump-log.sh \
  --files /var/lib/kafka-logs/my-topic-0/00000000000000000000.log \
  --verify-index-only

非正常重启后的检查清单

#!/bin/bash
# recovery-check.sh

BROKER_ID=1
LOG_DIR="/var/lib/kafka-logs"

echo "=== Recovery Check for Broker $BROKER_ID ==="

# 1. 检查日志目录是否存在且可写
if [ ! -w "$LOG_DIR" ]; then
  echo "[ERROR] Log directory not writable: $LOG_DIR"
  exit 1
fi

# 2. 检查是否有损坏的日志段
find $LOG_DIR -name "*.log" | while read logfile; do
  kafka-dump-log.sh --files "$logfile" --verify-index-only 2>&1 | grep -i error
done

# 3. 启动 Kafka(如有日志损坏,观察启动日志)
# 自动修复:unclean.shutdown.enable=true(谨慎使用)

7.2 使用 MirrorMaker 进行灾备恢复

MirrorMaker 是 Kafka 自带的数据镜像工具,可用于跨集群数据复制和灾难恢复。

配置 MirrorMaker

# consumer.properties - 源集群配置
cat > consumer.properties << 'CONSUMER'
bootstrap.servers=source-kafka:9092
group.id=mirrormaker-group
key.deserializer=org.apache.kafka.common.serialization.ByteArrayDeserializer
value.deserializer=org.apache.kafka.common.serialization.ByteArrayDeserializer
CONSUMER

# producer.properties - 目标集群配置
cat > producer.properties << 'PRODUCER'
bootstrap.servers=target-kafka:9092
key.serializer=org.apache.kafka.common.serialization.ByteArraySerializer
value.serializer=org.apache.kafka.common.serialization.ByteArraySerializer
PRODUCER

启动 MirrorMaker

kafka-mirror-maker.sh \
  --consumer.config consumer.properties \
  --producer.config producer.properties \
  --whitelist "orders|payments|logs-.*" \
  --num.streams 4

运维要点

  • --num.streams 控制并发消费线程数,建议与 Topic 分区数对齐
  • MirrorMaker 2.0(MM2)支持双向复制和偏移量同步,推荐在新版本中使用
  • 定期验证目标集群数据完整性

八、Prometheus + Grafana 面板

Prometheus + Grafana 是 Kafka 监控的主流方案。通过 JMX Exporter 采集 Kafka 的 JMX 指标,由 Prometheus 存储,Grafana 进行可视化展示。

8.1 部署 JMX Exporter

# 下载 JMX Exporter
wget https://repo1.maven.org/maven2/io/prometheus/jmx/jmx_prometheus_javaagent/0.19.0/jmx_prometheus_javaagent-0.19.0.jar -O /opt/jmx_exporter/javaagent.jar

# 创建 JMX Exporter 配置文件
mkdir -p /opt/jmx_exporter
cat > /opt/jmx_exporter/kafka.yml << 'JMXCONF'
lowercaseOutputName: true
lowercaseOutputLabelNames: true
rules:
  # UnderReplicatedPartitions
  - pattern: kafka.server<type=ReplicaManager, name=UnderReplicatedPartitions><>Value
    name: kafka_server_replicamanager_underreplicatedpartitions
    help: Number of under-replicated partitions
    type: GAUGE

  # OfflinePartitions
  - pattern: kafka.server<type=ReplicaManager, name=OfflinePartitions><>Value
    name: kafka_server_replicamanager_offlinepartitioncount
    help: Number of offline partitions
    type: GAUGE

  # ActiveControllerCount
  - pattern: kafka.controller<type=KafkaController, name=ActiveControllerCount><>Value
    name: kafka_controller_kafkacontroller_activecontrollercount
    help: Number of active controllers
    type: GAUGE

  # MessagesInPerSec
  - pattern: kafka.server<type=BrokerTopicMetrics, name=MessagesInPerSec><>OneMinuteRate
    name: kafka_server_brokertopicmetrics_messagesinpersec_rate
    help: Messages in per second
    type: GAUGE

  # BytesInPerSec / BytesOutPerSec
  - pattern: kafka.server<type=BrokerTopicMetrics, name=(BytesInPerSec|BytesOutPerSec)><>OneMinuteRate
    name: kafka_server_brokertopicmetrics_$1_rate
    help: Bytes in/out per second
    type: GAUGE

  # Request Latency
  - pattern: kafka.network<type=RequestMetrics, name=RequestsPerSec, request=(Produce|FetchConsumer|FetchFollower)><>OneMinuteRate
    name: kafka_network_requestmetrics_requestspersec
    labels:
      request: "$1"
    type: GAUGE
JMXCONF

8.2 配置 Kafka 启动参数

修改 server.propertieskafka-server-start.sh,添加 JMX Exporter Agent:

export KAFKA_OPTS="-javaagent:/opt/jmx_exporter/javaagent.jar=7071:/opt/jmx_exporter/kafka.yml"

启动后,访问 http://broker:7071/metrics 即可看到 Prometheus 格式的指标。

8.3 Prometheus 配置

# prometheus.yml
scrape_configs:
  - job_name: 'kafka'
    static_configs:
      - targets:
          - 'broker1:7071'
          - 'broker2:7071'
          - 'broker3:7071'
    relabel_configs:
      - source_labels: [__address__]
        target_label: instance

8.4 Grafana 面板配置

导入官方 Kafka Dashboard(ID: 721)或自建面板。以下是关键面板配置示例:

Broker 概览面板(Graph)

{
  "title": "Messages In Per Second",
  "targets": [
    {
      "expr": "kafka_server_brokertopicmetrics_messagesinpersec_rate",
      "legendFormat": "{{ instance }}"
    }
  ],
  "type": "timeseries",
  "fieldConfig": {
    "defaults": {
      "unit": "short"
    }
  }
}

Consumer Lag 面板(Table)

{
  "title": "Consumer Group Lag",
  "targets": [
    {
      "expr": "kafka_consumergroup_lag",
      "format": "table",
      "instant": true
    }
  ],
  "transformations": [
    {
      "id": "organize",
      "options": {
        "renameByName": {
          "group": "Consumer Group",
          "topic": "Topic",
          "partition": "Partition",
          "Value": "Lag"
        }
      }
    }
  ]
}

8.5 推荐的 Grafana Dashboard 组合

面板名称核心指标刷新频率
Cluster OverviewUnderReplicatedPartitions, ActiveControllerCount5s
Broker MetricsBytesIn/Out, MessagesIn, Request Latency10s
Consumer LagLag by Group/Topic/Partition15s
Disk UsageLog Size, Available Space1m
Topic DetailPartition Count, Replication Factor1m

九、升级策略与滚动重启

9.1 版本升级路径

Kafka 升级应遵循以下原则:

  1. 先升级 Zookeeper(如果使用较新版本需要)
  2. 滚动升级 Broker:逐台升级,确保集群始终可用
  3. 升级客户端:在所有 Broker 升级完成后,再升级生产者和消费者
# 查看当前版本
kafka-topics.sh --version

# 升级前备份配置
cp /opt/kafka/config/server.properties /opt/kafka/config/server.properties.bak

9.2 滚动重启流程

#!/bin/bash
# rolling-restart.sh

BROKERS=("broker1" "broker2" "broker3" "broker4" "broker5")

for broker in "${BROKERS[@]}"; do
  echo "=== Restarting $broker ==="

  # 1. 检查当前 ISR 状态
  echo "Checking ISR status before restart..."
  kafka-topics.sh --bootstrap-server $broker:9092 --describe | grep -i "isr"

  # 2. 优雅关闭 Broker
  echo "Gracefully stopping Kafka on $broker..."
  ssh $broker "kafka-server-stop.sh" || ssh $broker "systemctl stop kafka"

  # 3. 等待确认进程退出
  sleep 10
  ssh $broker "pgrep -f kafka.Kafka" && echo "WARNING: Kafka still running on $broker"

  # 4. 更新二进制文件(升级场景)
  # ssh $broker "tar -xzf kafka_2.13-3.6.0.tgz -C /opt/"
  # ssh $broker "ln -sfn /opt/kafka_2.13-3.6.0 /opt/kafka"

  # 5. 启动 Broker
  echo "Starting Kafka on $broker..."
  ssh $broker "systemctl start kafka"

  # 6. 验证 Broker 恢复
  sleep 30
  until kafka-broker-api-versions.sh --bootstrap-server $broker:9092 > /dev/null 2>&1; do
    echo "Waiting for $broker to be ready..."
    sleep 10
  done

  # 7. 验证 ISR 恢复
  echo "Checking ISR status after restart..."
  UNDER_REPLICATED=$(kafka-topics.sh --bootstrap-server $broker:9092 --describe | grep -c "UnderReplicated")
  echo "Under replicated partitions: $UNDER_REPLICATED"

  # 8. 等待集群稳定后再处理下一个
  echo "Waiting 60s for cluster stabilization..."
  sleep 60
  echo "$broker restart completed."
  echo ""
done

9.3 升级验证清单

检查项命令预期结果
Broker 版本kafka-broker-api-versions.sh显示新版本号
分区可用性kafka-topics.sh --describeLeader 分布正常,没有 UnderReplicated
消费正常kafka-consumer-groups.sh --describeLAG 稳定或递减
生产验证发送测试消息消息可正常发送和消费

十、常见故障排查

10.1 生产者发送失败

现象:生产者报错 NOT_ENOUGH_REPLICAS 或超时

# 检查分区 ISR 状态
kafka-topics.sh --bootstrap-server localhost:9092 --describe --topic orders

# 检查目标分区的 Leader 和 ISR
# 如果 ISR 数量小于 min.insync.replicas,生产消息会失败

# 解决方案:临时降低 acks 要求(不推荐长期)
# 或者修复副本同步问题

10.2 消费者不消费或消费缓慢

排查步骤

# 1. 查看消费者组状态
kafka-consumer-groups.sh \
  --bootstrap-server localhost:9092 \
  --describe --group my-group

# 2. 检查消费者是否还在组内
kafka-consumer-groups.sh \
  --bootstrap-server localhost:9092 \
  --describe --group my-group | grep -v "-"

# 3. 检查消费者的 Kafka 日志
journalctl -u kafka-consumer -f

# 4. 常见原因和解决
# - 消费者处理逻辑阻塞:优化业务代码
# - 分区分配不均:增加消费者实例
# - 数据倾斜:考虑重新分区策略

10.3 Broker OOM 或频繁 GC

# 查看 GC 日志
tail -f /var/log/kafka/kafkaServer-gc.log

# 检查堆内存使用
jstat -gcutil $(pgrep -f "kafka.Kafka") 1s 10

# 常见优化:
# 1. 增加堆内存
# export KAFKA_HEAP_OPTS="-Xms6g -Xmx6g"
# 2. 使用 G1GC 垃圾收集器
# export KAFKA_JVM_PERFORMANCE_OPTS="-server -XX:+UseG1GC"

10.4 磁盘 IO 瓶颈

# 检查磁盘 IO 等待
iostat -x 1

# 检查 Kafka 日志刷新配置
grep flush /opt/kafka/config/server.properties
# log.flush.interval.messages=10000
# log.flush.interval.ms=1000

# 优化建议
# 1. 使用独立磁盘存储 Kafka 日志
# 2. 考虑使用 RAID 10 平衡性能和可靠性
# 3. 调整刷新频率,避免过于频繁的 fsync

10.5 Zookeeper 连接异常

# 检查 Zookeeper 状态
zookeeper-shell.sh zk1:2181 <<< "stat"

# 检查 Broker 到 Zookeeper 的连接
netstat -an | grep 2181

# 检查 Zookeeper 会话超时配置
grep zookeeper.session.timeout.ms /opt/kafka/config/server.properties

十一、总结

Kafka 生产集群的运维是一项系统性工程,需要在监控、告警、容量管理和故障处理四个维度建立完整的体系。

监控层面:以 JMX 指标为核心,重点监控 UnderReplicatedPartitions、OfflinePartitions、ActiveControllerCount 和 ISRShrink 等指标。一个指标胜过千言万语,这些核心指标能够快速反映集群的健康状态。

Consumer Lag 层面:通过 kafka-consumer-groups.sh 进行日常巡查,结合 Burrow 实现自动化状态评估。Lag 不仅是一个数字,更是消费健康度的晴雨表 —— 持续增长的 Lag 预示着潜在的系统过载风险。

容量管理层面:Broker 的扩缩容和分区重分配需要谨慎规划。重分配操作涉及大量数据迁移,应在业务低峰期执行,并通过 throttle 参数限制对业务的影响。

数据恢复层面:熟悉日志恢复工具的使用,建立基于 MirrorMaker 的跨集群灾备方案。数据备份不是选项而是必需,定期的数据完整性验证同样不可或缺。

可视化层面:Prometheus + Grafana 的组合提供了直观的监控视图,合理设计的告警规则和 Dashboard 能大幅提升故障发现速度。

升级与运维层面:滚动重启保证升级期间的集群可用性,制定详细的升级验证清单可以避免遗漏关键检查项。

建立 Kafka 运维体系的核心原则是:可观测、可预警、可恢复。在生产环境中,任何一个维度的缺失都可能在关键时刻放大故障影响。希望本文提供的实战操作和配置示例,能帮助读者建立起适合自己业务场景的高可用 Kafka 运维方案。

继续阅读

探索更多技术文章

浏览归档,发现更多关于系统设计、工具链和工程实践的内容。

全部文章 返回首页

「kafka」更多文章

  1. 事件驱动架构:Event Sourcing、CQRS 与 Saga 模式
  2. Kafka 详解:分布式日志系统、ISR 与一致性保证
  3. Kafka 生产者与消费者实战:批量发送、ACK 策略与 Consumer Group