6.6 消息队列生产高可用——Kafka 集群部署与运维

预计阅读时间:21 分钟

📖 目录

mq_kafka01 介绍了 Kafka 的架构与基本操作。本文聚焦生产环境高可用:KRaft 3 节点集群部署、副本与 ISR 机制、Consumer Group 再平衡调优、Prometheus 监控及常见故障处理。

学习目标

  • 完成 KRaft 模式 3 节点 Kafka 集群的生产级部署
  • 深入理解副本与 ISR(In-Sync Replica)机制,掌握 min.insync.replicas 与 acks 的配合
  • 掌握 Consumer Group 再平衡的触发条件与调优方法,能解决消费堆积
  • 通过 JMX + Prometheus + Grafana 建立 Kafka 监控告警体系
  • 掌握常见故障(磁盘满、Leader 切换、分区丢失)的排查流程

前置知识

1. KRaft 模式 3 节点集群部署

Kafka 3.3+ 推荐 KRaft 替代 ZooKeeper,减少外部依赖。Controller 角色由 Broker 内置 Raft 协议选举产生。

# 每个节点 server.properties 核心配置
# node1
node.id=1
process.roles=broker,controller
listeners=PLAINTEXT://node1:9092,CONTROLLER://node1:9093
controller.quorum.voters=1@node1:9093,2@node2:9093,3@node3:9093
log.dirs=/data/kafka-logs
num.partitions=6
default.replication.factor=3
min.insync.replicas=2

# node2 / node3 同理,仅 node.id 和 listeners 不同
# 初始化与启动
KAFKA_CLUSTER_ID=$(bin/kafka-storage.sh random-uuid)
bin/kafka-storage.sh format -t $KAFKA_CLUSTER_ID -c config/kraft/server.properties
bin/kafka-server-start.sh -daemon config/kraft/server.properties

# 验证
bin/kafka-topics.sh --bootstrap-server node1:9092 --list
3 节点容错 Raft 需多数派存活,3 节点可容忍 1 节点宕机。生产建议 Controller 与 Broker 分离部署(3 Controller + N Broker)。

2. 副本与 ISR 机制

ISR(In-Sync Replicas)是高可用核心。只有与 Leader 保持同步的副本才属于 ISR 集合。

参数默认值推荐值说明
replication.factor13副本数
min.insync.replicas12最小确认副本数
acks1all等待所有 ISR 确认
unclean.leader.election.enabletruefalse禁止非 ISR 选 Leader
replica.lag.time.max.ms3000015000Follower 超时踢出 ISR
# 查看 ISR 状态
bin/kafka-topics.sh --bootstrap-server node1:9092 --describe --topic orders
# 输出:Partition: orders-0  Leader: 1  Replicas: 1,2,3  Isr: 1,2,3
# Isr 缺少节点说明该 Broker 落后或宕机
acks=all + min.insync.replicas=2 生产最低要求。ISR 数低于 min.insync.replicas 时 Producer 收到 NotEnoughReplicasException,写入失败。宁可拒绝写入也不丢消息。

3. Consumer Group 再平衡调优

再平衡期间所有消费者暂停消费,是消费延迟的主要来源。

触发场景说明影响
消费者加入/离开实例启动或下线全量重分配
心跳超时session.timeout.ms 超时误判为下线
处理超时max.poll.interval.ms 超时被踢出组
# 生产环境推荐配置
props.put("group.instance.id", "consumer-" + UUID.randomUUID());  // 静态成员
props.put("session.timeout.ms", "45000");
props.put("heartbeat.interval.ms", "15000");       // timeout/3
props.put("max.poll.interval.ms", "300000");       // 单次处理最大5分钟
props.put("partition.assignment.strategy",
    "org.apache.kafka.clients.consumer.CooperativeStickyAssignor");
静态成员 设置 group.instance.id 后重启不触发再平衡,滚动升级可零停机。CooperativeStickyAssignor(Kafka 2.4+)实现增量再平衡,只迁移必要分区,显著缩短耗时。

4. 监控指标——JMX + Prometheus

Kafka 通过 JMX 暴露运行指标,配合 JMX Exporter + Prometheus + Grafana 实现全栈可观测。

# 启动 Broker 时挂载 JMX Exporter agent
KAFKA_OPTS="-javaagent:/opt/jmx_prometheus_javaagent.jar=7071:/opt/kafka-jmx.yml" \
  bin/kafka-server-start.sh -daemon config/kraft/server.properties
指标说明告警阈值
UnderReplicatedPartitionsISR 不足的分区数>0 持续 5 分钟
ActiveControllerCount活跃 Controller 数!= 1
RequestHandlerAvgIdlePercent请求处理空闲率<30%
NetworkProcessorAvgIdlePercent网络线程空闲率<30%
MessagesInPerSec每秒写入消息数按基线调整
Consumer Lag消费者组堆积量>10000 持续增长
# Prometheus 告警规则示例
groups:
  - name: kafka_alerts
    rules:
      - alert: KafkaUnderReplicatedPartitions
        expr: kafka_server_ReplicaManager_UnderReplicatedPartitions > 0
        for: 5m
        labels: { severity: critical }
        annotations:
          summary: "Kafka ISR 副本不足: {{ $value }} 个分区"
      - alert: KafkaConsumerLagHigh
        expr: kafka_consumergroup_lag_sum > 10000
        for: 10m
        labels: { severity: warning }
        annotations:
          summary: "消费者组 {{ $labels.consumergroup }} 堆积: {{ $value }}"

4.1 关键指标告警阈值

告警不是"有指标就报",要区分故障信号与噪音。下表给出生产集群核心指标的阈值建议,告警级别用于区分通知渠道(critical 走电话/钉钉,warning 走邮件/群消息)。

指标含义建议阈值级别
UnderReplicatedPartitionsISR 不足的分区数>0 持续 5 分钟critical
OfflinePartitionsCount无 Leader 的分区数>0 立即告警critical
Consumer Lag 总和消费组未消费消息数>10000 且 10 分钟趋势上升warning
ActiveControllerCount活跃 Controller 数!= 1critical
Controller 切换次数1 小时内 Controller 变更次数>2 次/小时warning
磁盘使用率log.dirs 所在分区使用率≥80% warning,≥90% criticalwarning/critical
网络吞吐Broker 网卡出口流量≥ 网卡带宽 70% 持续 10 分钟warning
RequestHandlerAvgIdlePercent请求处理线程空闲率<30%critical

4.2 告警规则表达式

# ISR 缩容与离线分区
expr: kafka_server_ReplicaManager_UnderReplicatedPartitions > 0
expr: kafka_server_KafkaServer_OfflinePartitionsCount > 0

# Controller 抖动(1 小时内切换超过 2 次说明集群不稳定,常见于节点网络抖动)
expr: changes(kafka_controller_KafkaController_ActiveControllerCount[1h]) > 2

# 消费组堆积且持续增长(用 delta 判断趋势,避免偶发尖峰误报)
expr: kafka_consumergroup_lag_sum > 10000 and
      delta(kafka_consumergroup_lag_sum[10m]) > 0

# 磁盘使用率(node_exporter 采集,按数据盘挂载点聚合)
expr: (1 - node_filesystem_avail_bytes{fstype=~"ext4|xfs"}
      / node_filesystem_size_bytes{fstype=~"ext4|xfs"}) * 100 >= 80

# 网卡吞吐达到带宽 70%(30G 网卡示例,bit/s 换算后比较)
expr: rate(node_network_transmit_bytes_total{device="bond0"}[5m]) * 8
      / 30000000000 > 0.7
阈值要随业务调 上表是通用起点。峰值业务(大促)的 Lag 阈值应放宽 5-10 倍;磁盘告警要与 retention 清理窗口错开,避免"清理前必然告警"的永久噪音。先让告警跑两周,按误报率收敛阈值。

Kafka 与 RabbitMQ 高可用方案对比

维度Kafka (KRaft)RabbitMQ (镜像队列)
消息模型发布-订阅(消费者组)点对点(队列)
数据持久化日志追加写入磁盘内存 + 磁盘(可配置)
消费语义at-least-once / exactly-onceat-least-once
高可用机制副本 + ISR + Leader 选举镜像队列 + Quorum 队列
消息回溯支持(按 offset / 时间戳)不支持
吞吐量极高(百万级 msg/s)高(万级 msg/s)
适用场景日志流、事件溯源、流处理任务队列、RPC、延迟消息
选型建议 需要消息回溯、高吞吐、流处理能力时选 Kafka;需要复杂路由(topic exchange)、延迟消息、优先级队列时选 RabbitMQ。两者可以组合使用:Kafka 负责日志流,RabbitMQ 负责业务任务队列。

5. 常见故障处理

5.1 磁盘满

# 查看磁盘占用
du -sh /data/kafka-logs/*/ | sort -rh | head -20

# 临时调整保留大小
bin/kafka-configs.sh --bootstrap-server node1:9092 \
  --entity-type topics --entity-name huge-topic \
  --alter --add-config retention.bytes=10737418240  # 10GB

# 紧急清理(需先停 Broker)
bin/kafka-server-stop.sh
rm -rf /data/kafka-logs/huge-topic-0/*.log
bin/kafka-server-start.sh -daemon config/kraft/server.properties

5.2 消费堆积

# 查看堆积
bin/kafka-consumer-groups.sh --bootstrap-server node1:9092 \
  --group order-service --describe

# 扩容消费者(不超过分区数)
# 调大 max.poll.records 增加单次拉取量
# 临时增加分区数(注意 key 哈希会变化)
bin/kafka-topics.sh --bootstrap-server node1:9092 \
  --alter --topic orders --partitions 12

5.3 Broker 宕机

# 确认宕机节点
bin/kafka-topics.sh --bootstrap-server node2:9092 --describe | grep "Leader: -1"

# Kafka 自动将 Leader 切换到其他 ISR 副本
# 恢复宕机节点
bin/kafka-server-start.sh -daemon config/kraft/server.properties
# 等待 Replica Fetcher 追上,ISR 自动扩展
# 监控 UnderReplicatedPartitions 归零即恢复完成
故障类型影响恢复时间优先级
Broker 宕机Leader Partition 不可用自动秒级切换P0
磁盘满写入被拒绝手动清理 5-30 分钟P0
消费堆积消费延迟增大扩容后分钟级P1
ISR 缩减数据冗余降低网络恢复后自动P1

6. 生产高可用检查清单

检查项要求验证方式
集群节点数≥ 3 节点kafka-metadata-quorum.sh
副本因子≥ 3kafka-topics.sh --describe
ISR 最小值min.insync.replicas ≥ 2kafka-configs.sh
acks 配置Producer acks=all客户端配置
非 ISR 选举unclean.leader.election.enable=falsekafka-configs.sh
监控告警UnderReplicatedPartitions 告警Prometheus 告警规则
磁盘监控使用率 <80%node_exporter + 告警
Consumer Lag<10000Prometheus / consumer-groups.sh

7. 集群滚动重启与版本升级

滚动重启(Rolling Restart)指逐台停止、升级、启动 Broker,全程不中断读写。任何生产变更(内核参数、JVM 参数、版本升级、安全补丁)都应走滚动重启,而不是整集群停机。

7.1 Kafka 滚动重启完整步骤

# 步骤 0:前置检查——只有集群健康才能开始
bin/kafka-topics.sh --bootstrap-server node1:9092 --describe --topic orders \
  | grep -c "Isr:"    # 对比分区总数,确认没有 ISR 缩减
bin/kafka-metadata-quorum.sh --bootstrap-server node1:9092 describe --status
# ActiveControllerCount 应为 1

# 步骤 1:开启优雅关闭(server.properties,重启后生效)
controlled.shutdown.enable=true
controlled.shutdown.max.retries=3
controlled.shutdown.retry.backoff.ms=5000

# 步骤 2:优雅停止第一台 Broker(Leader 自动迁走,等待 ISR 重配置完成)
bin/kafka-server-stop.sh node1    # 或 kill -TERM $(cat /data/kafka.pid)
tail -f /data/kafka-logs/server.log | grep -E "Shutdown|shutdown"

# 步骤 3:确认分区 Leader 已切换到存活节点(不应出现 Leader: -1)
bin/kafka-topics.sh --bootstrap-server node2:9092 --describe --topic orders \
  | awk '{print $2, $4}'

# 步骤 4:在窗口内完成该节点操作(打补丁/改配置/换内核)
# 步骤 5:启动 Broker,等待其追上最新数据并回到 ISR
bin/kafka-server-start.sh -daemon config/kraft/server.properties

# 步骤 6:确认 UnderReplicatedPartitions 归零后,再重启下一台
# 步骤 7:全部完成后执行 Preferred Leader 选举,恢复副本平衡
bin/kafka-leader-election.sh --bootstrap-server node1:9092 \
  --election-type preferred --all-topic-partitions
Controller 节点注意 KRaft 模式下 Controller 与 Broker 同进程,逐台重启时 3 节点仲裁始终有 2 台存活,安全。但绝不允许两台 Controller 同时重启——会丢失仲裁导致元数据操作失败;必须等上一台启动并加入仲裁后再操作下一台。

7.2 版本升级策略:先 broker 后 client

Kafka 官方支持跨一个次要版本升级(如 3.5 → 3.6),且新 Broker 向后兼容旧客户端。因此标准顺序是先升级 Broker,再升级客户端:新客户端可能使用旧 Broker 不认识的请求协议(分区发现、配额协商),先升 Client 会引入"客户端能发、服务端不认"的风险窗口;而旧客户端请求新 Broker 永远兼容。

# 1. 检查当前版本与协议版本
bin/kafka-broker-api-versions.sh --bootstrap-server node1:9092

# 2. 逐台滚动升级 Broker 二进制(过程同 7.1,注意升级后配置兼容性)
# 3. 全部 Broker 升级完成后,再升级集群元数据版本
bin/kafka-storage.sh upgrade -c config/kraft/server.properties

# 4. 观察 1-2 天运行稳定后,再升级客户端库
#    Producer/Consumer 依赖:kafka-clients 3.9.x -> 4.0.x
# 5. 生产验证期保留一个旧版本客户端做冒烟测试,确认兼容后再全部替换

7.3 RabbitMQ 集群滚动重启

RabbitMQ 多节点集群的滚动重启核心是 stop_app(只停应用、不停 Erlang 节点),节点优雅退出集群,其余节点自动接管队列与镜像。

# 前置检查:确认集群无分区(Partitions 列表为空)
rabbitmqctl cluster_status | grep -A10 "Partitions"

# 逐台操作:优雅停止应用(不动 Erlang 虚拟机)
rabbitmqctl stop_app
# 升级 Erlang/插件或打补丁
# 重新启动应用,节点自动重新加入集群并同步镜像
rabbitmqctl start_app

# 等待镜像队列同步完成(synchronised_slave_pids 数量 = 镜像副本数)
rabbitmqctl list_queues name slave_pids synchronised_slave_pids
# 确认 running_nodes 恢复 3 台后再操作下一台
rabbitmqctl cluster_status

8. 分区扩容与再均衡

8.1 增加分区数

分区数是 Topic 并发度的上限:一个分区只能被组内一个消费者消费,分区不足时消费能力被锁死。但扩容的副作用必须提前评估:

  • key 路由变化hash(key) % 分区数 的分布彻底改变,同 key 消息会落到不同分区,按 key 保序的场景(订单流水按订单号聚合)被破坏,且无法回退
  • 已有数据不迁移:新分区为空,老分区数据不动,流量只会逐步分散到新分区
  • 分区数只增不减:Kafka 不支持减少分区数,误扩容只能重建 Topic 并迁移数据
# 扩容示例:6 -> 12 分区(适用于无 key 的日志类 Topic)
bin/kafka-topics.sh --bootstrap-server node1:9092 \
  --alter --topic audit-log --partitions 12
# 验证
bin/kafka-topics.sh --bootstrap-server node1:9092 --describe --topic audit-log
带 key 的 Topic 扩容等于重排 例如订单 Topic 按 order_id 分区,扩容后同一订单的消息分散到不同分区,消费端聚合逻辑错乱。此类 Topic 扩容前必须与业务方确认,或改在消息体内携带全局序号做补偿排序。

8.2 调整副本因子

副本因子不足(如误建 replication.factor=1)可以原地补齐,无需重建 Topic。用 reassignment JSON 描述每个分区期望的副本列表:

# 1. 生成期望副本分配(topics.json 声明要处理的 Topic)
cat > topics.json <<EOF
{"topics":[{"topic":"orders"}],"version":1}
EOF
bin/kafka-reassign-partitions.sh --bootstrap-server node1:9092 \
  --generate --topics-to-move-json-file topics.json \
  --broker-list 1,2,3 > reassignment.json

# 2. 手动编辑 reassignment.json,把每个分区的 replicas 补到 3 副本
#    {"topic":"orders","partition":0,"replicas":[1,2,3]}

# 3. 执行并限速 50MB/s,避免数据搬迁挤占业务带宽
bin/kafka-reassign-partitions.sh --bootstrap-server node1:9092 \
  --reassignment-json-file reassignment.json --execute --throttle 50000000

# 4. 轮询验证(未完成前输出 "in progress")
bin/kafka-reassign-partitions.sh --bootstrap-server node1:9092 \
  --reassignment-json-file reassignment.json --verify

# 5. 完成后解除限速(重要!否则限速永久生效,集群吞吐被压制)
bin/kafka-reassign-partitions.sh --bootstrap-server node1:9092 \
  --reassignment-json-file reassignment.json --execute --throttle -1

8.3 分区均衡(kafka-reassign-partitions 完整流程)

新节点加入、旧节点退役、或长期运行后分区分布不均(各 Broker 磁盘使用率差异大),用官方工具做全集群均衡。标准流程:generate 生成方案 -> 审阅 -> execute 执行 -> verify 验证。

# 1. 声明要移动的 Topic(空列表表示全集群)
cat > topics.json <<EOF
{"topics":[{"topic":"orders"},{"topic":"audit-log"}],"version":1}
EOF

# 2. 生成均衡方案(输出可执行的 reassignment.json)
bin/kafka-reassign-partitions.sh --bootstrap-server node1:9092 \
  --generate --topics-to-move-json-file topics.json --broker-list 1,2,3 \
  > reassignment.json

# 3. 人工审阅:确认每个分区 3 副本、Leader 分散在不同 Broker
# 4. 低峰期执行,限速 100MB/s
bin/kafka-reassign-partitions.sh --bootstrap-server node1:9092 \
  --reassignment-json-file reassignment.json --execute --throttle 100000000

# 5. 持续观察执行进度(期间集群承担额外复制流量,注意网络指标)
bin/kafka-reassign-partitions.sh --bootstrap-server node1:9092 \
  --reassignment-json-file reassignment.json --verify

# 6. 解除限速并确认各 Broker 磁盘水位接近
bin/kafka-reassign-partitions.sh --bootstrap-server node1:9092 \
  --reassignment-json-file reassignment.json --execute --throttle -1
限速参数 --throttle 单位是字节/秒,按节点维度汇总。写小了大数据搬迁拖到数小时,写大了挤占业务带宽。建议从峰值带宽的 30% 起步,观察复制速率与业务时延后调整。

9. 数据可靠性——acks、ISR 与 HW 截断

9.1 acks=all 与 min.insync.replicas 的配合

这两个参数必须成对理解:acks=all 要求 Leader 等所有 ISR 成员确认后才返回成功;min.insync.replicas 定义"至少有几个 ISR 成员"才能接受写入。前者决定"等谁",后者决定"最少等几个"。

副本因子min.insync.replicas可容忍宕机数行为
321任意 1 台宕机写入正常,数据至少 2 副本
3301 台宕机即拒绝写入(NotEnoughReplicas),可用性差
532推荐生产配置:容忍 2 台宕机且写入不丢
110无冗余,任何宕机都可能丢数据,仅测试用

9.2 HW(High Watermark)与数据截断

HW 是所有副本已同步的位点,消费者只能读到 HW 之前的数据。Leader 崩溃后,新 Leader 会把本地数据截断到自身 HW:已被 acks=all 确认、但尚未复制到全部 ISR 的消息,在极端场景(多个副本同时宕机)下仍可能丢失。

# 查看分区 HW 与 LEO(Log End Offset)
bin/kafka-get-offsets.sh \
  --broker-list node1:9092 --topic orders --time -1   # 最新 offset(LEO)
bin/kafka-get-offsets.sh \
  --broker-list node1:9092 --topic orders --time -2   # 最旧 offset

# 截断发生的痕迹:Broker 日志出现
# WARN [ReplicaManager broker=2] ... local log end offset ... below high watermark ...
# 说明该副本被截断(丢掉了未同步的分段),需人工核对数据影响

9.3 数据丢失场景与防护

场景发生条件防护措施
非 ISR 副本上位unclean.leader.election.enable=true 且 Leader 全挂设为 false,宁可短暂不可用也不丢数据
acks=1 写入丢失Producer 只等 Leader 确认,Leader 复制前崩溃acks=all + min.insync.replicas=2
HW 截断丢数据Follower 长期落后,新 Leader 按自身 HW 截断缩短 replica.lag.time.max.ms、监控 ISR 状态
单盘故障磁盘损坏,log.dirs 数据不可读replication.factor=3、JBOD 多目录冗余
机房级故障整机房断电/网络中断跨机房 MirrorMaker2 容灾(异步复制)

一句话总结可靠性模型:acks=all + min.insync.replicas=2 + replication.factor=3 + 关闭 unclean 选举,在单台 Broker 宕机下做到零丢失;跨机房容灾只能依赖异步复制,接受秒级窗口的潜在丢失。

10. 故障演练

高可用配置是否真的生效,只有把故障真实注入一次才知道。演练必须选低峰期、有回滚预案、并记录"预期行为 vs 实际行为"。

10.1 演练场景矩阵

场景注入方法预期行为恢复步骤
进程被杀kill -9 <broker_pid>30-60 秒内 Leader 切换,acks=all 且 ISR≥2 时写入不中断重启 Broker,ISR 自动追平
拔网线/断网iptables -A INPUT -s node2 -j DROPreplica.lag.time.max.ms 后该节点掉出 ISR,不再影响写入恢复网络,ISR 自动回补
磁盘满dd if=/dev/zero of=/data/kafka-logs/fill bs=1M count=10000该 Broker 写入失败;分区若只剩一个 ISR 则拒绝写入清理 Segment 或扩容,恢复后 ISR 追平
Controller 单点故障kill Controller 角色进程其余节点重新选举 Controller,元数据操作短暂不可用重启该节点,自动回归
时钟偏移date -s "+300 seconds"(先停 NTP)心跳时间戳错乱,节点被误判掉线校准时钟并恢复 NTP,ISR 回补

10.2 演练实例:kill -9 一个 Broker

# 演练前:记录基线(ISR 全绿、Lag 归零)
bin/kafka-topics.sh --bootstrap-server node1:9092 --describe --topic orders \
  | grep Isr | sort > /tmp/baseline.txt

# 注入故障:强杀 node3(无优雅关闭)
ssh node3 'kill -9 $(pgrep -f kafka.Kafka)'

# 观察(30 秒内应看到):
# 1. 其他 Broker 日志出现 Leader 迁移:Elected leader ... for partition
# 2. Prometheus 中 UnderReplicatedPartitions 短暂 >0 后归零
# 3. 生产流量无报错(前提:acks=all 且 ISR≥2)

# 恢复:重启 node3,等待 ISR 追平(落后越多恢复越慢)
ssh node3 'bin/kafka-server-start.sh -daemon config/kraft/server.properties'
watch -n 10 'bin/kafka-topics.sh --bootstrap-server node1:9092 --describe --topic orders | grep node3'

# 演练后:校验基线一致、执行 Preferred Leader 选举、写演练报告
diff /tmp/baseline.txt <(bin/kafka-topics.sh --bootstrap-server node1:9092 \
  --describe --topic orders | grep Isr | sort)

10.3 演练注意事项

  • 演练前通知业务方与值班人员,明确恢复窗口与"叫停"信号(谁有权中止演练)
  • 一次只注入一个故障,观察完整恢复链路后再进行下一个,避免故障叠加掩盖根因
  • 演练必须有停止条件:如超时未自动恢复、业务报错率超阈值,立即执行恢复预案
  • kill -9 强杀与优雅关闭(kill -TERM)都要演练,两者恢复路径不同
  • 结果纳入 SOP:每季度一次,新成员必须独立完成一份"故障报告"

常见错误

问题现象排查与解决
NotEnoughReplicasException生产者写入失败,提示 ISR 数不足检查 min.insync.replicas 设置,确认所有 ISR 副本存活,运行 kafka-topics.sh --describe 查看 Isr 列表
UnderReplicatedPartitions 告警Prometheus 持续告警 ISR 不足检查宕机或慢副本的 Broker 状态,确认网络连通性,等待 Replica Fetcher 追上后 Isr 自动恢复
再平衡耗时过长消费者组再平衡期间消费长时间暂停启用 CooperativeStickyAssignor 增量再平衡,设置 group.instance.id 静态成员避免不必要再平衡
磁盘写入拒绝Broker 日志报磁盘满,生产被拒绝紧急清理过期 Segment,调整 log.retention.bytesretention.ms,扩容磁盘
Controller 选举失败集群无法正常进行 Leader 选举确认 controller.quorum.voters 配置正确,检查 Controller 节点间网络连通性,确保多数派存活
扩容后 key 顺序错乱同 key 消息被不同消费者消费,业务聚合结果错误带 key 的 Topic 扩容前评估路由影响,优先考虑分区均衡替代扩容;无法避免时与业务方确认补偿方案
reassign 后集群长期低吞吐数据搬迁持续占用带宽,业务时延升高执行 reassignment 时设置 --throttle 限速,完成后务必用 --throttle -1 解除限速
滚动重启导致 Controller 双宕两台 Controller 同时重启,集群元数据服务失联严格逐台重启,等上一台加入仲裁(quorum)后再操作下一台

最佳实践

建议说明
3 节点起步 + min.insync.replicas=23 节点容忍 1 节点宕机,ISR 最小值设为 2 保证写入至少 2 副本确认
关闭 unclean.leader.election禁止非 ISR 副本竞选 Leader,宁可短暂不可用也不丢数据
启用 CooperativeStickyAssignor增量再平衡只迁移必要分区,大幅缩短再平衡耗时,降低消费暂停影响
配置 JMX + Prometheus + Grafana全栈可观测,重点关注 UnderReplicatedPartitions、Consumer Lag、RequestHandlerIdle 三项指标
定期执行高可用检查清单对照第 6 节检查清单逐项验证,纳入运维 SOP,确保生产环境持续高可用
滚动重启前先检查集群健康UnderReplicatedPartitions 归零、ActiveControllerCount=1 才能开始,否则故障与变更混在一起无法定位
定期做故障演练kill -9、断网、磁盘满至少每季度一次,验证监控告警与自动恢复链路真实可用

故障排查案例:ISR 缩减导致写入失败

现象:Producer 持续报 NotEnoughReplicasException,Prometheus 告警 UnderReplicatedPartitions > 0 持续 10 分钟。

排查:执行 kafka-topics.sh --describe --topic orders | grep Isr,发现某分区 Isr 列表缺少 node3,而 node3 上的 kafka-server.log 出现 Replication Lag Exceeded 警告。

根因:node3 的磁盘 I/O 被日志写入高峰占满,Replica Fetcher 无法及时同步 Leader 数据,导致落后时间超过 replica.lag.time.max.ms(30s)阈值,被踢出 ISR。

修复:① 将 Kafka 数据盘与系统盘物理分离;② 在 server.properties 中增加 num.io.threads=16 提升 I/O 并发度;③ 重启 node3 后观察 ISR 自动恢复。

消息积压应急处理模板

# 步骤 1:确认堆积规模与趋势
bin/kafka-consumer-groups.sh --bootstrap-server node1:9092 \
  --group order-service --describe | awk '{sum+=$6} END{print "Total lag:", sum}'

# 步骤 2:临时增加消费者实例(不超过分区数)
# 如果分区数不足,先扩容分区(注意 key 哈希变化)
bin/kafka-topics.sh --bootstrap-server node1:9092 \
  --alter --topic orders --partitions 12

# 步骤 3:调大单次拉取量
# 在消费者配置中增加
# max.poll.records=500
# fetch.min.bytes=1048576

# 步骤 4:监控消费进度恢复
watch -n 5 'bin/kafka-consumer-groups.sh --bootstrap-server node1:9092 \
  --group order-service --describe | awk "{print \$1, \$2, \$6}"'
应急原则 消费堆积的根因通常是消费者处理能力不足或下游依赖(数据库/API)变慢。扩容只是缓解手段,必须同步排查下游瓶颈。

练习题

  1. 部署一个 3 节点 KRaft 集群,设置 replication.factor=3min.insync.replicas=2acks=all。创建 Topic 并模拟 1 个 Broker 宕机,观察 Leader 自动切换过程,验证写入不中断。
  2. 配置 Prometheus 监控 Kafka JMX 指标,编写告警规则:当 UnderReplicatedPartitions > 0 持续 5 分钟触发 P0 告警,当 Consumer Lag > 10000 持续 10 分钟触发 P1 告警。用 Grafana 仪表板可视化。
  3. 对同一 Topic 分别使用 RangeAssignor(默认)和 CooperativeStickyAssignor 启动消费者组,模拟消费者加入/离开,对比两种策略的再平衡耗时和消费暂停时长。
  4. 对 3 节点集群执行一次完整滚动重启(含 1 台 Broker 优雅关闭与重启),全程记录:Leader 切换耗时、ISR 恢复耗时、业务是否中断,并对比 kill -9 强杀与优雅关闭的恢复差异。
  5. 模拟 1 台 Broker 磁盘满:观察写入失败与 ISR 变化,练习清理过期 Segment 恢复集群的完整流程,并验证 Prometheus 磁盘告警是否按预期触发。

学习检查点

学完本章后,请检验自己是否掌握以下内容:

检查项自测问题验证方法
概念理解能用自己的话解释消息队列高可用的主从同步和仲裁机制尝试向他人讲解
命令操作能不查文档完成 RabbitMQ 镜像队列或 Kafka 集群的高可用配置在终端实际执行
原理掌握能说出消息队列的持久化、复制和故障恢复原理画出流程图
故障排查能独立排查消息队列集群脑裂或数据不一致的问题模拟故障并修复
最佳实践能说明为什么需要为消息队列配置监控和告警对比不同方案

本章总结

Kafka 高可用的本质是"冗余 + 多数派":副本数保证数据冗余,ISR 与 acks=all 保证写入不丢,Leader 选举保证故障自动切换。生产环境要把这三件事钉死:副本数 ≥3、min.insync.replicas=2、监控先行。没有监控的高可用集群只是"看起来可用"——故障发生时才发现不可用,是运维最大的事故。

延伸阅读

↑ 回到顶部