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 切换、分区丢失)的排查流程
前置知识
- mq_kafka01 Kafka 架构与基本操作——本文是其生产实战延伸
- 6.4:消息队列入门 消息队列入门——消息队列基础概念
- 3.12:系统监控与告警 系统监控与告警——Prometheus/Grafana 的使用基础
- 6.15:HA 集群 HA 集群——Raft 共识与多节点部署的基本思想
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
2. 副本与 ISR 机制
ISR(In-Sync Replicas)是高可用核心。只有与 Leader 保持同步的副本才属于 ISR 集合。
| 参数 | 默认值 | 推荐值 | 说明 |
|---|---|---|---|
replication.factor | 1 | 3 | 副本数 |
min.insync.replicas | 1 | 2 | 最小确认副本数 |
acks | 1 | all | 等待所有 ISR 确认 |
unclean.leader.election.enable | true | false | 禁止非 ISR 选 Leader |
replica.lag.time.max.ms | 30000 | 15000 | Follower 超时踢出 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 落后或宕机
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
| 指标 | 说明 | 告警阈值 |
|---|---|---|
| UnderReplicatedPartitions | ISR 不足的分区数 | >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 走邮件/群消息)。
| 指标 | 含义 | 建议阈值 | 级别 |
|---|---|---|---|
| UnderReplicatedPartitions | ISR 不足的分区数 | >0 持续 5 分钟 | critical |
| OfflinePartitionsCount | 无 Leader 的分区数 | >0 立即告警 | critical |
| Consumer Lag 总和 | 消费组未消费消息数 | >10000 且 10 分钟趋势上升 | warning |
| ActiveControllerCount | 活跃 Controller 数 | != 1 | critical |
| Controller 切换次数 | 1 小时内 Controller 变更次数 | >2 次/小时 | warning |
| 磁盘使用率 | log.dirs 所在分区使用率 | ≥80% warning,≥90% critical | warning/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
Kafka 与 RabbitMQ 高可用方案对比
| 维度 | Kafka (KRaft) | RabbitMQ (镜像队列) |
|---|---|---|
| 消息模型 | 发布-订阅(消费者组) | 点对点(队列) |
| 数据持久化 | 日志追加写入磁盘 | 内存 + 磁盘(可配置) |
| 消费语义 | at-least-once / exactly-once | at-least-once |
| 高可用机制 | 副本 + ISR + Leader 选举 | 镜像队列 + Quorum 队列 |
| 消息回溯 | 支持(按 offset / 时间戳) | 不支持 |
| 吞吐量 | 极高(百万级 msg/s) | 高(万级 msg/s) |
| 适用场景 | 日志流、事件溯源、流处理 | 任务队列、RPC、延迟消息 |
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 |
| 副本因子 | ≥ 3 | kafka-topics.sh --describe |
| ISR 最小值 | min.insync.replicas ≥ 2 | kafka-configs.sh |
| acks 配置 | Producer acks=all | 客户端配置 |
| 非 ISR 选举 | unclean.leader.election.enable=false | kafka-configs.sh |
| 监控告警 | UnderReplicatedPartitions 告警 | Prometheus 告警规则 |
| 磁盘监控 | 使用率 <80% | node_exporter + 告警 |
| Consumer Lag | <10000 | Prometheus / 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
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
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 | 可容忍宕机数 | 行为 |
|---|---|---|---|
| 3 | 2 | 1 | 任意 1 台宕机写入正常,数据至少 2 副本 |
| 3 | 3 | 0 | 1 台宕机即拒绝写入(NotEnoughReplicas),可用性差 |
| 5 | 3 | 2 | 推荐生产配置:容忍 2 台宕机且写入不丢 |
| 1 | 1 | 0 | 无冗余,任何宕机都可能丢数据,仅测试用 |
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 DROP | replica.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.bytes 或 retention.ms,扩容磁盘 |
| Controller 选举失败 | 集群无法正常进行 Leader 选举 | 确认 controller.quorum.voters 配置正确,检查 Controller 节点间网络连通性,确保多数派存活 |
| 扩容后 key 顺序错乱 | 同 key 消息被不同消费者消费,业务聚合结果错误 | 带 key 的 Topic 扩容前评估路由影响,优先考虑分区均衡替代扩容;无法避免时与业务方确认补偿方案 |
| reassign 后集群长期低吞吐 | 数据搬迁持续占用带宽,业务时延升高 | 执行 reassignment 时设置 --throttle 限速,完成后务必用 --throttle -1 解除限速 |
| 滚动重启导致 Controller 双宕 | 两台 Controller 同时重启,集群元数据服务失联 | 严格逐台重启,等上一台加入仲裁(quorum)后再操作下一台 |
最佳实践
| 建议 | 说明 |
|---|---|
| 3 节点起步 + min.insync.replicas=2 | 3 节点容忍 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}"'
练习题
- 部署一个 3 节点 KRaft 集群,设置
replication.factor=3、min.insync.replicas=2、acks=all。创建 Topic 并模拟 1 个 Broker 宕机,观察 Leader 自动切换过程,验证写入不中断。 - 配置 Prometheus 监控 Kafka JMX 指标,编写告警规则:当 UnderReplicatedPartitions > 0 持续 5 分钟触发 P0 告警,当 Consumer Lag > 10000 持续 10 分钟触发 P1 告警。用 Grafana 仪表板可视化。
- 对同一 Topic 分别使用
RangeAssignor(默认)和CooperativeStickyAssignor启动消费者组,模拟消费者加入/离开,对比两种策略的再平衡耗时和消费暂停时长。 - 对 3 节点集群执行一次完整滚动重启(含 1 台 Broker 优雅关闭与重启),全程记录:Leader 切换耗时、ISR 恢复耗时、业务是否中断,并对比 kill -9 强杀与优雅关闭的恢复差异。
- 模拟 1 台 Broker 磁盘满:观察写入失败与 ISR 变化,练习清理过期 Segment 恢复集群的完整流程,并验证 Prometheus 磁盘告警是否按预期触发。
学习检查点
学完本章后,请检验自己是否掌握以下内容:
| 检查项 | 自测问题 | 验证方法 |
|---|---|---|
| 概念理解 | 能用自己的话解释消息队列高可用的主从同步和仲裁机制 | 尝试向他人讲解 |
| 命令操作 | 能不查文档完成 RabbitMQ 镜像队列或 Kafka 集群的高可用配置 | 在终端实际执行 |
| 原理掌握 | 能说出消息队列的持久化、复制和故障恢复原理 | 画出流程图 |
| 故障排查 | 能独立排查消息队列集群脑裂或数据不一致的问题 | 模拟故障并修复 |
| 最佳实践 | 能说明为什么需要为消息队列配置监控和告警 | 对比不同方案 |
本章总结
Kafka 高可用的本质是"冗余 + 多数派":副本数保证数据冗余,ISR 与 acks=all 保证写入不丢,Leader 选举保证故障自动切换。生产环境要把这三件事钉死:副本数 ≥3、min.insync.replicas=2、监控先行。没有监控的高可用集群只是"看起来可用"——故障发生时才发现不可用,是运维最大的事故。
延伸阅读
- mq_kafka01 Apache Kafka 消息队列深入——架构、安装与生产调优
- 3.12:系统监控与告警 系统监控与告警(Prometheus + Grafana 全栈监控)
- 6.4:消息队列入门 消息队列入门——RabbitMQ 与 Redis Stream
- 4.7:压力测试实战 压力测试实战——用 wrk/hey 测试消息系统吞吐