Skip to content

第 10 章 Controller 与 KRaft:Kafka 的「大脑改造手术」

目标读者:被「Kafka 不是要去掉 ZooKeeper 了吗?」「KRaft 跟 Raft 是一回事吗?」「我们集群还在用 ZK,要不要升级?」这些问题困扰的同学。

学完你会:能讲清楚 Kafka 元数据管理的两代架构、Controller 的角色、KRaft 解决了 ZK 时代哪些痛点;遇到「Active Controller 频繁切换」「集群启动慢」这类告警,知道该看哪个日志、动哪个配置;做生产架构选型时,能给出「新集群用 KRaft,老集群什么时候迁移」的合理判断。


0. 导读:Kafka 也有「大脑」

绝大多数 Kafka 教程都把 Broker 当成对等的劳动者,事实上集群里有一个「特殊角色」:Controller

  • 数据面(Data Plane):Producer 写、Consumer 读、Broker 之间复制——这些是分布式日志的「血管」;
  • 控制面(Control Plane):谁是 Leader、谁在 ISR、Topic 有哪几个分区、配置变更如何下发——这些是分布式系统的「神经」。

本章不讲数据面(前面 9 章已经讲透),专门讲控制面:谁来管这些元数据?变更怎么传播?发生故障怎么恢复?

Kafka 的控制面经历了一次根本性的架构重写:从 0.x 到 2.7 的 ZooKeeper 时代,到 2.8 引入的 KRaft(Kafka Raft)预览,到 3.3 KRaft 生产可用,到 4.0 完全删除 ZK。这是一次堪比 MySQL 把元数据从 MyISAM 表搬到 InnoDB 的「内脏移植」。


1. ZK 时代:Controller 是个「值班长」

1.1 为什么需要 Controller

集群里有很多 Broker,每个 Broker 上有很多分区。为了让所有 Broker 看到一致的元数据视图,必须有人负责:

  • Topic 创建/删除/扩分区时,决定新分区分配到哪几个 Broker
  • Broker 上下线时,重新选 Leader、调整 ISR;
  • 把这些变更广播给所有 Broker。

Kafka 选择了一个常见但有点偷懒的方案:从所有 Broker 中选一个出来当「值班长」,承担上述决策与广播职责。这个值班长就叫 Controller

注意:Controller 不是单独部署的进程,它就是某一个普通 Broker,多了一些职责而已。

1.2 ZK 时代的 Controller 选举

Kafka 利用 ZooKeeper 的「临时节点 + 监听器」实现了一个非常朴素的选举:

所有 Broker 启动后争抢创建 ZK 上的 /controller 临时节点
  谁先创建成功 → 谁就是 Controller,并把自己的 broker_id 写进去
其它 Broker 监听 /controller 节点:
  一旦该节点消失(Controller 进程死了,临时节点被 ZK 自动删除)
  所有 Broker 重新争抢 → 选出新 Controller

Controller 的核心职责(ZK 时代):

  1. 监听 ZK 上 /brokers/ids/*/brokers/topics/*/admin/* 等节点,感知集群变化;
  2. 计算新的 Leader / ISR / 分区分配;
  3. 通过 UpdateMetadataRequest / LeaderAndIsrRequest / StopReplicaRequest 等 RPC,把变更广播给所有 Broker;
  4. 把元数据写回 ZK(持久化)。

1.3 ZK 元数据布局

/controller                           (临时节点:当前 Controller 是谁)
/controller_epoch                     (单调递增:每次切 Controller 就 +1)
/brokers/
  ├── ids/
  │   ├── 1   (Broker1 心跳)
  │   ├── 2   (Broker2 心跳)
  │   └── 3   (Broker3 心跳)
  └── topics/
      ├── learn.09.orders/
      │   ├── partitions/0/state    (Leader, ISR, controller_epoch, leader_epoch)
      │   ├── partitions/1/state
      │   └── partitions/2/state
      └── ...
/admin/
  ├── delete_topics/
  ├── reassign_partitions/
  └── preferred_replica_election/
/config/
  ├── topics/learn.09.orders/        (Topic 级配置覆盖)
  ├── brokers/1/                     (Broker 级动态配置)
  └── users/...
/consumers/                          (旧版 Consumer offsets,0.9 之后已废弃)

每一个变更都是「先写 ZK,再发 RPC」,整个集群的「真理」都在 ZK 里。

1.4 生活类比:ZK 是个「公司组织部」

把 Kafka 集群比作一家公司。Broker 是各部门员工,Controller 是当前的轮值总经理,ZK 是「组织部」——所有人的入职信息、岗位安排都在它的档案柜里。

  • 谁离职、谁晋升、谁调岗,都先去组织部改档案,再由轮值总经理通知所有部门;
  • 总经理离职,靠组织部的「值班登记表」(临时节点)通知大家「谁来接班」;
  • 所有员工每天打卡都要往组织部签到(心跳)。

2. ZK 方案的历史包袱

ZK 帮 Kafka 跑了 10 年,但随着集群规模越做越大,问题越来越明显。

2.1 元数据变更顺序问题

ZK 不保证「变更顺序」与「通知顺序」严格一致。Controller 也不能保证 RPC 到达所有 Broker 的顺序一致。曾经出现过:

  • Broker A 先收到「Topic X 分区 0 Leader = B1」;
  • Broker A 接着收到「Topic X 分区 0 Leader = B2」;
  • 但同样一对消息到达 Broker B 的时候顺序颠倒了,B 上看到的 Leader 是 B1。

为了避免这种乱序,Kafka 不得不在每个 RPC 里带 controller_epoch + leader_epoch,Broker 收到旧 epoch 的指令直接丢弃。这套补丁能跑,但所有逻辑都被 epoch 比较污染,代码非常复杂。

2.2 规模上限:10 万分区是块铁顶

ZK 时代的元数据有几个硬性瓶颈:

  • ZK 单节点容量:所有 Topic / Partition 元数据都堆在 ZK 内存,超过 50 万 znode 后 GC 时间会让 ZK 卡顿;
  • Controller 启动时拉取元数据:新 Controller 上线第一件事,是把 ZK 上几万个分区的元数据全部拉下来,构建本地内存视图——这个操作可能要好几分钟。
  • 元数据广播 O(N):分区数 P,Broker 数 B,一次集群变更要发 O(B*P) 量级的 RPC。变更频繁时网络打满。

社区当年公开提到过:单集群 10 万分区 是 ZK 时代实际使用的天花板。LinkedIn / Confluent 内部的大集群早就被这个问题逼到「拆集群」「禁用某些功能」的境地。

2.3 Controller 故障恢复慢

ZK 时代 Controller 死了,新 Controller 上来要做:

  1. 从 ZK 把所有 Topic / Partition 元数据全部拉下来(O(P));
  2. 从 ZK 拉所有 Broker 状态(O(B));
  3. 比对自己的本地缓存,重新计算每个分区是否需要重新选 Leader;
  4. 给所有 Broker 发一遍 UpdateMetadataRequest

这个「重启重建」过程在大集群上可能要 几十秒到几分钟——这段时间里 Topic 创建会卡住、分区 Leader 切换不能进行。可用性受影响。

2.4 「两套真理」的运维负担

集群的真相分散在两个系统里:Kafka 自己的日志 + ZK 的 znode。

  • 安装、监控、备份要对接两套系统;
  • 安全/认证要对接两套(ZK ACL + Kafka SASL);
  • 升级、扩容要协调两套;
  • ZK 自己也是分布式系统,也有自己的故障模式(脑裂、leader 选举抖动),出问题时 Kafka 同样躺枪。

每一个 Kafka 运维都被 ZK 至少咬过一次。


3. KRaft 详解:把元数据搬进 Kafka 自己

3.1 KIP-500 的核心思路

KRaft(Kafka Raft)= 「把元数据本身存成一个 Kafka Topic」+「用 Raft 协议管理这个 Topic 的 Leader 与一致性」

具体说:

  • 创建一个内部 Topic:__cluster_metadata,1 个分区、N 副本(N = Controller Quorum 大小,通常 3 或 5);
  • 这个 Topic 的副本不分散在所有 Broker 上,而只放在指定的 Controller 节点
  • Controller 们之间用 Raft 协议选举出 Active Controller(Raft Leader),其它是 Voter Follower;
  • 集群所有元数据变更(建 Topic、改 ISR、Broker 上下线……)都变成一条事件记录追加到 __cluster_metadata 里;
  • Broker 通过订阅 __cluster_metadata 来获取元数据更新——它不再被「推」,而是「拉」

一句话:Kafka 用 Kafka 管 Kafka

3.2 KRaft Quorum 与 Active Controller 选举

集群里有 3 ~ 5 个 Controller 节点,组成一个 Raft Quorum:
┌─────────────┐    ┌─────────────┐    ┌─────────────┐
│ Controller1 │ ←→ │ Controller2 │ ←→ │ Controller3 │
│  (Leader)   │    │  (Follower) │    │  (Follower) │
└─────────────┘    └─────────────┘    └─────────────┘
           ↓ 都在维护同一个 __cluster_metadata 日志

Raft 协议的核心:任何写入都要被多数派 Quorum 持久化才算 commit。3 节点要 2 个,5 节点要 3 个。

Active Controller 就是 Raft Leader,它独占处理元数据写入;其他节点只 fetch 复制日志、参与投票。如果 Active 挂了,剩余节点通过选举(带 term 的投票)选出新 Leader。

KRaft 的 Raft 实现并非完全照搬经典论文,而是为了配合 Kafka 现有传输层做了简化(KIP-595 / 631 / 642)。它叫 KRaft 的原因正是「Kafka-flavored Raft」。

3.3 元数据事件溯源(Event Sourcing)

ZK 时代是「最终态存储」:znode 里就是元数据的当前样子。 KRaft 时代是「事件溯源」:__cluster_metadata 里存的是一系列变更事件

offset=0    RegisterBrokerRecord(broker=1, ...)
offset=1    RegisterBrokerRecord(broker=2, ...)
offset=2    TopicRecord(topic=foo)
offset=3    PartitionRecord(topic=foo, part=0, leader=1, isr=[1,2,3], leader_epoch=0)
offset=4    PartitionRecord(topic=foo, part=1, leader=2, isr=[1,2,3], leader_epoch=0)
offset=5    PartitionChangeRecord(topic=foo, part=0, new_leader=2, leader_epoch=1)
offset=6    UnregisterBrokerRecord(broker=1)
...

每一个 record 是一个不可变事件,Active Controller 把它追加到日志、Quorum 多数派持久化、commit。Broker 端通过 fetch 这些 record,在内存里重放得到最新元数据视图。

好处

  • 强顺序:日志的 offset 自然定义了变更顺序,不再需要 epoch 比较;
  • 天然审计:所有历史变更都在日志里,回放可以重建任何时刻的集群状态;
  • 快速同步:Broker 启动时只要拉到最新 offset,就完成元数据初始化(不用一个个 znode 遍历)。

3.4 元数据快照

事件日志会越来越长,重启时全量重放会变慢。Raft 协议用 快照(Snapshot) 解决这个问题:

  • Active Controller 周期性地把当前元数据状态序列化成快照文件,写入磁盘 + 同步到 Quorum;
  • 老的事件日志可以截断(保留快照之后的);
  • 新节点上线时先加载快照,再追加增量事件。

KRaft 快照:

  • 默认每 1 小时或日志增长 100 MB 触发;
  • 文件名 00000000000000000123-0000000000.checkpoint(offset + epoch);
  • 配合 KIP-630 实现增量快照,节省 IO。

3.5 Broker 如何获取元数据

ZK 时代:Controller 给所有 Broker。 KRaft 时代:Broker __cluster_metadata

  • 每个 Broker 启动时,向 Active Controller 注册 RegisterBrokerRequest
  • 之后 Broker 不断发 FetchRequest 拉取 __cluster_metadata 的新事件(就跟普通 Consumer 拉日志一样);
  • 同时 Broker 通过 BrokerHeartbeatRequest 周期向 Controller 上报自己的状态。

这种「拉模式」让 Broker 上下线、元数据变更都不需要 Controller 主动重发,元数据传播变得幂等且可恢复——Broker 短暂失联,恢复后只要补拉错过的 offset 就行。


4. ZK vs KRaft 启动恢复对比

4.1 ZK 时代:O(N) 全量拉取

新 Controller 起来:
  for each topic in zk:
    for each partition in topic:
      读 /brokers/topics/.../partitions/X/state
  for each broker in zk:
    读 /brokers/ids/X
  ===> 几万分区时这一步要几分钟
  
  对所有分区重新计算 Leader / ISR
  广播 UpdateMetadataRequest 给每个 Broker

10 万分区集群,整个过程可能 30 秒到几分钟。

4.2 KRaft 时代:O(1) 加载快照 + 增量

新节点起来:
  loadSnapshot()                  ← 一次磁盘读
  fetchAndApply(增量 records)     ← 从最近快照之后开始追
  ===> 通常 < 1 秒
  
没有「重新计算 + 广播」的步骤——所有元数据都已经在 __cluster_metadata 里了

Confluent 在 3.3 release blog 里给出过实测:100 万分区集群下,Active Controller 故障切换恢复时间从 ZK 时代的 5+ 分钟 降到 KRaft 的 < 30 秒

4.3 元数据广播延迟对比

                         ZK 时代                    KRaft 时代
Topic 创建可见时延        Controller 推给所有 Broker  Broker 自主拉取(毫秒级)
                         (O(B) RPC)                
分区 Leader 切换          Controller 单独发           包含在元数据日志里
                         LeaderAndIsr 给相关 Broker   一并 fetch
Controller 故障切换       几十秒到几分钟              < 1 秒(Raft 选举)
集群可承载分区上限         实测约 10 万                100 万 +

5. 从 ZK 迁移到 KRaft:双写路径

Kafka 3.4 ~ 3.7 提供了滚动迁移方案,业务不需要停服。这不是一次性切换,而是一个多步骤过程。

5.1 阶段总览

┌─────────────┐  step 1   ┌─────────────────────┐  step 2   ┌─────────────────────┐  step 3   ┌─────────────┐
│ 纯 ZK 模式  │ ────────→ │ KRaft Migration 模式│ ────────→ │ KRaft 模式(带 ZK) │ ────────→ │ 纯 KRaft   │
│ Broker only │           │ Controller 双写    │           │ ZK 旁路             │           │             │
└─────────────┘           └─────────────────────┘           └─────────────────────┘           └─────────────┘

5.2 详细步骤(基于 Kafka 3.6+)

前置:所有 Broker 升级到 Kafka 3.6+,配置 inter.broker.protocol.version=3.6 等元数据版本。

Step 1:起一个新的 KRaft Controller Quorum(3 节点),配置 kafka.controller.zookeeper.metadata.migration.enable=true,并指向旧 ZK:

properties
process.roles=controller
node.id=3000
controller.quorum.voters=3000@new-controller-1:9093,3001@new-controller-2:9093,3002@new-controller-3:9093
controller.listener.names=CONTROLLER
zookeeper.metadata.migration.enable=true
zookeeper.connect=old-zk:2181

Step 2:新 Controller Quorum 起来后,Active Controller 会从 ZK 把所有元数据拷贝__cluster_metadata,这时进入 dual-write(双写) 模式:

  • 新元数据写入到 __cluster_metadata AND 同时写回 ZK;
  • 旧 Broker 仍从 ZK 读元数据,新 Broker 从 KRaft 读;
  • 集群处在「ZK 是真相 ←→ KRaft 是真相」的过渡态。

Step 3:滚动重启每个 Broker,把它们改成 KRaft 模式(process.roles=brokercontroller.quorum.voters=...、不再配置 zookeeper.connect):

properties
process.roles=broker
node.id=1
controller.quorum.voters=3000@new-controller-1:9093,3001@new-controller-2:9093,3002@new-controller-3:9093
# 不再配置 zookeeper.connect

Step 4:所有 Broker 都切到 KRaft 后,把 Active Controller 配置里的 zookeeper.metadata.migration.enable 改成 false,并去掉 ZK 连接信息。这一步之后 ZK 完全没用了。

Step 5:可以销毁 ZK 集群。

5.3 迁移前必读的「禁飞清单」

Apache 官方在 migration 文档 列出了几个不支持迁移的功能,迁移前必须确认你没用到:

  • 旧的 kafka-acls.sh --authorizer kafka.security.auth.SimpleAclAuthorizer(ZK 版 ACL)—— KRaft 用 kafka.security.authorizer.StandardAuthorizer
  • 旧版 SCRAM 用户存储在 ZK 中——需要先迁移到 kafka-storage.sh 的本地存储;
  • 任何直接读写 ZK 的第三方工具(比如某些老 Kafka Manager);
  • 多 ZK chroot(zookeeper.connect=host:2181/chroot1,/chroot2)。

5.4 Kafka 4.0 完全移除 ZK

Kafka 4.0(2025 年发布)已经彻底删除 ZooKeeper 相关代码:

  • zookeeper.connect 配置直接报错;
  • zkCli.sh 工具下架;
  • 所有迁移工作必须在 3.x 完成。

也就是说,如果你现在还在用 ZK 模式,升级到 4.x 之前必须先迁到 KRaft——这是硬性要求,没有「跳过迁移直接升 4.x」的选项。


6. KRaft 的关键配置

6.1 角色与拓扑

KRaft 模式下,Kafka 进程通过 process.roles 决定自己是什么:

配置角色说明
process.roles=broker纯 Broker(接客户端读写)不参与元数据投票
process.roles=controller纯 Controller(管元数据)通常 3 或 5 个,组成 Quorum
process.roles=broker,controller兼任(合并模式)学习/小集群用,生产不推荐

生产推荐:Controller 与 Broker 分开部署。Controller 节点专门跑 Raft,避免被业务流量挤压;3 ~ 5 台够用,不需要随集群规模线性扩。

6.2 必填配置示例

Controller 节点process.roles=controller):

properties
process.roles=controller
node.id=1
controller.quorum.voters=1@controller1:9093,2@controller2:9093,3@controller3:9093
controller.listener.names=CONTROLLER
listeners=CONTROLLER://0.0.0.0:9093
log.dirs=/var/lib/kafka/metadata

Broker 节点process.roles=broker):

properties
process.roles=broker
node.id=101
controller.quorum.voters=1@controller1:9093,2@controller2:9093,3@controller3:9093
controller.listener.names=CONTROLLER
listeners=PLAINTEXT://0.0.0.0:9092
advertised.listeners=PLAINTEXT://broker-1.example.com:9092
log.dirs=/var/lib/kafka/data

关键点解读

  • controller.quorum.voters 三节点都填完全一样的值,包括端口;
  • controller.listener.names=CONTROLLER 告诉 Kafka 哪个监听器是 Quorum 内部通信用的;
  • node.id 全局唯一,不能重复;
  • 第一次启动前必须用 kafka-storage.sh format 初始化 metadata.log.dir(生成 meta.properties、cluster ID 等)。

6.3 第一次启动的 storage format

KRaft 模式启动前必须先 format(ZK 时代不需要):

bash
# 生成一个集群 UUID(22 字符 base64)
CLUSTER_ID=$(/opt/kafka/bin/kafka-storage.sh random-uuid)
echo $CLUSTER_ID

# 在每个节点上 format
/opt/kafka/bin/kafka-storage.sh format \
  --config /opt/kafka/config/kraft/server.properties \
  --cluster-id $CLUSTER_ID

format 会在 log.dirs 下生成 meta.properties,里面记着 cluster_id 和 node_id。所有节点的 cluster_id 必须相同


7. 实操:观察 KRaft 集群

本教程的 docker-compose.yml 已经是 KRaft 模式,跑起来后可以:

7.1 看集群基本信息

bash
docker compose exec kafka1 \
  /opt/kafka/bin/kafka-broker-api-versions.sh \
  --bootstrap-server localhost:19092 | head -20

7.2 看 Controller 是谁

bash
docker compose exec kafka1 \
  /opt/kafka/bin/kafka-metadata-quorum.sh \
  --bootstrap-server localhost:19092 \
  describe --status

输出示例:

ClusterId:              MkU3OEVBNTcwNTJENDM2Qk
LeaderId:               1
LeaderEpoch:            5
HighWatermark:          1234
MaxFollowerLag:         0
MaxFollowerLagTimeMs:   0
CurrentVoters:          [1, 2, 3]
CurrentObservers:       []

LeaderId=1 就是当前 Active Controller

7.3 看每个 Voter 的复制状态

bash
docker compose exec kafka1 \
  /opt/kafka/bin/kafka-metadata-quorum.sh \
  --bootstrap-server localhost:19092 \
  describe --replication

输出示例:

NodeId  LogEndOffset  Lag  LastFetchTimestamp  LastCaughtUpTimestamp  Status
1       1234          0    -                   -                      Leader
2       1234          0    1700000000000       1700000000000          Follower
3       1233          1    1700000000000       1700000000000          Follower

Lag 是这个 Follower 落后 Leader 多少 record;Status 区分 Leader / Follower / Observer。

7.4 解析元数据日志

KRaft 时代有个新工具 kafka-metadata-shell.sh,能读 __cluster_metadata 日志:

bash
docker compose exec kafka1 \
  /opt/kafka/bin/kafka-metadata-shell.sh \
  --snapshot /var/lib/kafka/data/__cluster_metadata-0/00000000000000000000.log

进入交互式 shell,可以 lscat 元数据节点:

>> ls /
brokers  topics  configs  acls  ...

>> cat /topics/learn.09.orders/0/data
{
  "partitionId" : 0,
  "topicId" : "...",
  "replicas" : [ 1, 2, 3 ],
  "isr" : [ 1, 2, 3 ],
  "leader" : 1,
  "leaderEpoch" : 0,
  "partitionEpoch" : 0
}

非常方便排查「Topic 元数据是不是和 Broker 一致」「ACL 是否生效」之类问题。


8. KRaft 的「Observer」概念(容易混淆)

回顾第 9 章里提到过:开源 Kafka 的数据副本只有 Leader / Follower,没有 Observer。但 KRaft 内部 有 Observer 的概念,含义完全不同:

KRaft Observer = 不参与 Raft 投票,但会同步 __cluster_metadata 的节点

典型用法:

  • Broker 同步元数据:每个 Broker 都是 __cluster_metadata 的 Observer——它拉日志但不投票;
  • 新 Voter 加入前的 catch-up:Raft 协议要求新 Voter 加入 Quorum 前必须先 catch up;这段过渡期它就是 Observer;
  • 企业版的「投票冷备」:某些场景下放一个不投票的副本节点,纯做监听 / 灾备读。

不要把它和「数据副本 Observer」(Confluent MRC 商业特性)搞混。


9. KRaft 的「踩坑实录」(社区常见问题)

9.1 Quorum 配置漂移

3 个节点的 controller.quorum.voters 必须完全一致,包括 host:port、id。如果某节点配错了一个,会出现「能选出 Leader 但 Broker 拉元数据失败」「不停 reconnect」的诡异现象。

排查

bash
docker compose exec kafka1 cat /opt/kafka/config/kraft/server.properties \
  | grep controller.quorum.voters

三个节点比对,必须字符级别相同。

9.2 cluster.id 不一致

每个节点的 meta.propertiescluster.id 必须相同。如果你在某台机器上 format 时用了不同的 --cluster-id,启动会报 InconsistentClusterIdException

修复:把那台节点的 log.dirs 全删掉,重新 format。

9.3 Quorum 节点数选 4 不选 3

Quorum 容错数 = floor((N-1)/2)

Voter 数 N容错数说明
10单节点,挂就完蛋(学习用)
31推荐起步
41比 3 还差!多一个机器但容错没多
52大集群推荐
73极少需要

为什么 4 比 3 还差?因为 Raft 多数派 = floor(N/2)+1:N=4 时多数派是 3,需要至少 3 个节点活着;N=3 时多数派也是 2,需要至少 2 个节点活着。两者都只能容忍 1 节点宕,但 4 节点要多养一台机器,毫无收益。永远选奇数

9.4 __cluster_metadata 占盘异常

如果集群元数据变更频繁、快照间隔配置过大,__cluster_metadata 日志会涨得很大。可以调:

properties
metadata.log.max.snapshot.interval.ms=3600000   # 1 小时一次快照
metadata.max.idle.interval.ms=500
metadata.log.segment.bytes=1048576              # 1MB 切段

9.5 Active Controller 频繁切换

通常是 Quorum 节点之间网络抖动导致心跳失败、触发选举。排查方向:

  • controller.quorum.election.timeout.ms(默认 1000)适当调大;
  • 检查 Quorum 节点之间网络延迟(理想 < 10ms);
  • 检查 Controller 节点 GC(如果它兼任 Broker,业务流量会拉长 GC pause)。

10. 与其他系统的对比

                      ┌────────────────────────┬────────────────────────┐
                      │ 元数据存储             │ 选主协议              │
┌──────────────────┬──┴────────────────────────┼────────────────────────┤
│ Kafka (ZK 时代)   │ ZooKeeper               │ /controller 临时节点    │
│ Kafka (KRaft)    │ Kafka 自己 (__cluster..) │ Raft (KRaft 变种)       │
│ Pulsar           │ ZooKeeper + BookKeeper   │ ZK 临时节点             │
│ RabbitMQ Quorum  │ Raft 日志                │ Raft                    │
│ RocketMQ DLedger │ Raft 日志                │ Raft                    │
│ TiDB / TiKV      │ etcd (Raft)              │ Raft                    │
│ ClickHouse Keeper│ Raft (clickhouse-keeper) │ Raft                    │
│ Etcd / Consul    │ 自己 (Raft)              │ Raft                    │
└──────────────────┴───────────────────────────┴────────────────────────┘

可以看到一个明显趋势:新一代分布式系统都用 Raft 自管元数据,不再依赖外部 ZK。Kafka 走的是同一条路,只是历史包袱比新系统重得多。


11. 常用命令速查

任务命令
看 Quorum 状态kafka-metadata-quorum.sh ... describe --status
看 Quorum 复制延迟kafka-metadata-quorum.sh ... describe --replication
解析元数据日志kafka-metadata-shell.sh --snapshot ...
看 Active Controllerkafka-metadata-quorum.sh ... describe --status 看 LeaderId
生成 cluster.idkafka-storage.sh random-uuid
第一次 formatkafka-storage.sh format --config server.properties --cluster-id ...
列出当前 Brokerkafka-broker-api-versions.sh --bootstrap-server ...
看集群 ID / ControllerPython AdminClient list_topics().controller_id / .cluster_id

12. 小结:KRaft 是 Kafka 的「灵魂改造」

把本章压缩成 5 点:

  1. Controller 是集群的「值班长」,负责所有元数据决策与广播——ZK 时代是某个 Broker 兼任,KRaft 时代是独立的 Quorum。
  2. ZK 时代的痛点:变更顺序问题、10 万分区上限、Controller 切换慢、两套真理的运维负担。
  3. KRaft 的核心:把元数据当成 Kafka 自己的 Topic(__cluster_metadata),用 Raft 保证一致性,Broker 主动 fetch(拉)而不是 Controller 推。
  4. 迁移路径:3.6+ 的双写迁移,4.0 完全删 ZK——还在 ZK 上的集群必须在 4.x 升级前完成迁移。
  5. 生产建议:新集群直接 KRaft,Controller 与 Broker 分离部署,Quorum 选 3 / 5(不要选 4)。

13. 面试高频题

Q1:Kafka 为什么要去 ZK?KRaft 解决了什么核心问题?

考察点:架构演进的理解。

标准答案

  • 「去 ZK」不是因为 ZK 不好,而是 Kafka 不想再维护两套分布式系统:ZK 是「外部依赖」,运维、安全、监控、备份都要单独搞一套。
  • KRaft 解决的核心问题:
    1. 元数据规模:从 ZK 时代的 10 万分区上限,提升到 100 万以上;
    2. Controller 故障切换速度:从分钟级降到秒级;
    3. 元数据顺序一致性:用 __cluster_metadata 的 offset 顺序天然定义变更顺序,不再依赖 epoch 比较补丁;
    4. 运维简化:只剩一套系统,监控、安全、ACL 全部统一在 Kafka 里。

加分项:能补充「KRaft 的 Raft 是为 Kafka 量身定制的,不是照搬经典 Raft 论文(KIP-595/631)」。

Q2:KRaft 和 ZK 时代的 Controller 选举有什么不同?

考察点:选举协议本质。

标准答案

  • ZK 时代:所有 Broker 争抢 /controller 临时节点,谁先创建谁当选;ZK 用 ZAB 协议保证临时节点唯一性。Controller 就是某个普通 Broker 兼任。
  • KRaft 时代:3~5 个独立的 Controller 节点组成 Raft Quorum,通过标准 Raft 选举(带 term 的投票)选出 Active Controller;其它节点是 Voter Follower。
  • 关键差异
    • ZK 时代选举依赖外部协调系统;KRaft 自己实现共识协议;
    • ZK 时代 Controller 故障要等 ZK session timeout(默认 6s+)才感知;KRaft 是 Raft 内部心跳,毫秒级;
    • ZK 时代变更必须「写 ZK」+「广播给所有 Broker」,两步可能不一致;KRaft 是「commit 到 Raft 日志」一步到位,Broker 自主拉取。

Q3:KRaft 的 __cluster_metadata 是什么?它和普通 Topic 有什么区别?

考察点:对元数据存储的细节理解。

标准答案

  • __cluster_metadata 是 KRaft 时代专门存「集群元数据事件」的 内部 Topic,1 分区,副本数等于 Controller Quorum 大小。
  • 与普通 Topic 的区别:
    1. 副本只放在 Controller 节点上,不分散到所有 Broker;
    2. 复制协议是 Raft,要 Quorum 多数派持久化才 commit;普通 Topic 是 ISR 协议;
    3. 写入方限定 Active Controller,普通 Topic 任何 Producer 都能写;
    4. 每条 record 是结构化的元数据事件(如 PartitionRecordTopicRecordRegisterBrokerRecord),有 Schema;
    5. 配合快照机制减少日志体积。
  • Broker 通过 fetch 这个 Topic 的事件,在内存里重放得到最新的集群元数据视图。

加分项:能讲到「PartitionRecord 包含 partitionEpoch,跟 leaderEpoch 一起解决一致性」。

Q4:从 ZK 迁移到 KRaft 怎么做?停服吗?

考察点:实操经验。

标准答案

  • 不需要停服,是滚动迁移。基于 Kafka 3.6+ 的 dual-write 模式:
    1. 升级所有 Broker 到 3.6+;
    2. 起一组新的 KRaft Controller Quorum,开 zookeeper.metadata.migration.enable=true 并指向旧 ZK;
    3. 新 Controller 自动从 ZK 拷元数据到 __cluster_metadata,进入双写模式;
    4. 滚动重启每个 Broker,改成 KRaft 模式,去掉 zookeeper.connect;
    5. 全部切完后关闭 dual-write,销毁 ZK;
  • 注意事项:旧 ACL(SimpleAclAuthorizer)、SCRAM 用户、第三方 ZK 工具都要先迁移;
  • Kafka 4.0 完全删除 ZK,必须在升级到 4.x 前完成迁移。

Q5:KRaft 模式下 Controller Quorum 选 3 节点还是 5 节点?为什么不选 4?

考察点:Raft 容错原理。

标准答案

  • Raft 容错数 = floor((N-1)/2),多数派 = floor(N/2)+1
    • 3 节点:容忍 1 宕,多数派 = 2;
    • 4 节点:容忍 1 宕,多数派 = 3(比 3 节点更差!要 3 个活着才行);
    • 5 节点:容忍 2 宕,多数派 = 3;
  • 偶数节点不增加容错能力,只增加投票延迟和成本;
  • 必须选奇数:3 是入门,5 是大集群推荐,7 极少需要。

Q6:Kafka 4.0 之前我还在用 ZK,能直接升级到 4.0 吗?

考察点:升级路径。

标准答案

  • 不能。Kafka 4.0 已经完全删除 ZooKeeper 相关代码,包括 zookeeper.connect 配置、ZK 元数据迁移代码、zkCli.sh 工具。
  • 必须先在 3.x 完成迁移到 KRaft,再升级到 4.x。推荐路径:
    1. 升级到 3.6 或更高版本;
    2. 按 KIP-866 文档完成 ZK → KRaft 迁移;
    3. 验证集群在 KRaft 模式下稳定运行一段时间;
    4. 升级到 4.x。
  • 如果业务允许停服,也可以选择「拉一个新 KRaft 集群 + MirrorMaker 2.0 镜像数据」的方案,但这个对 Topic 命名、Offset 迁移要求高。

Q7:KRaft 模式下,Controller 节点能不能同时当 Broker?

考察点:部署模式选择。

标准答案

  • 技术上可以,配置 process.roles=broker,controller(合并模式 / Combined Mode)。
  • 学习环境、小集群(< 10 Broker)可以用,节省机器;
  • 生产环境强烈不推荐
    1. Controller Raft 需要稳定低延迟,被业务流量挤压会触发频繁选举;
    2. 业务流量打满磁盘 / 网络时,元数据同步会滞后;
    3. Controller 节点的 OS-level 调优(如 IO 优先级、网络 buffer)需求和 Broker 不同;
  • 生产推荐:3 ~ 5 台专属 Controller(小机器即可,CPU 4 核内存 8G 足够)+ N 台 Broker(高配机器),物理隔离。

Q8:「Active Controller 频繁切换」如何排查?

考察点:实战排障。

标准答案

  • 频繁切换说明 Quorum 之间心跳/选举出现问题。排查方向:
    1. 网络抖动:Quorum 节点之间网络延迟应稳定 < 10ms;用 ping / mtr 验证;
    2. GC 长停顿:Controller JVM GC pause 超过 controller.quorum.election.timeout.ms(默认 1s)会被认为失联;用 -Xlog:gc* 看 GC 日志;
    3. 磁盘 IO 瓶颈__cluster_metadata 写入慢会拖慢 Raft commit;iostat 看是不是磁盘忙;
    4. 配置不一致:三个节点 controller.quorum.voters / cluster.id 必须完全相同;
    5. 节点同时跑 Broker:合并模式下业务流量挤压 Controller,建议拆分;
  • 临时缓解:调大 controller.quorum.election.timeout.ms 到 2~5 秒;
  • 监控指标:JMX kafka.controller:type=KafkaController,name=ActiveControllerCount(应当稳定为 1)。

下一章我们将聚焦客户端最常踩坑的领域——消费者组与 Rebalance:JoinGroup / SyncGroup / Heartbeat 协议、4 种分区分配策略、Eager 与 Cooperative 的世纪对决、Static Membership 如何治抖、Lag 排查全流程。

🎬 可视化演示

演示加载缓慢或样式异常?点此在新标签页打开 ↗

💻 示例代码

bash
#!/usr/bin/env bash
# =============================================================================
# dump_metadata_log.sh
#   一键解析 KRaft 集群的 __cluster_metadata 日志,把里头的元数据事件
#   (TopicRecord、PartitionRecord、PartitionChangeRecord、RegisterBrokerRecord …)
#   人类可读地打印出来。
#
# 这是排查「Topic 创建后 Broker 看不到」「ACL 不生效」「Quorum 元数据撕裂」
# 类问题的「核武器」工具。
#
# 用法(在 docker-compose 起来的 KRaft 集群里):
#   ./dump_metadata_log.sh                        # 默认 dump kafka1 的元数据
#   ./dump_metadata_log.sh kafka2                 # dump kafka2
#   ./dump_metadata_log.sh kafka1 shell           # 进入交互式 metadata-shell
#   ./dump_metadata_log.sh kafka1 latest          # 只看最新 segment 的可读输出
#
# 依赖:
#   * Kafka 3.3+(kafka-metadata-shell.sh 和 kafka-dump-log.sh 都自带)
#   * docker compose(容器名 kafka1/kafka2/kafka3)
#
# 关键命令解释:
#   1) kafka-metadata-shell.sh    交互式查询元数据快照(支持 ls / cat / find)
#   2) kafka-dump-log.sh          按 record 一条条打印日志原文
#   3) kafka-metadata-quorum.sh   看 Quorum 状态 / 各 Voter 复制位
# =============================================================================

set -euo pipefail

CONTAINER=${1:-kafka1}
MODE=${2:-dump}
META_DIR="/var/lib/kafka/data/__cluster_metadata-0"

usage() {
  cat <<EOF
用法:$0 [container] [mode]
  container : kafka1 / kafka2 / kafka3 (默认 kafka1)
  mode      : dump      → 全量打印所有 segment(默认)
              latest    → 只打印当前最新 segment
              shell     → 进入交互式 kafka-metadata-shell.sh
              quorum    → 显示 Quorum 状态(不解析日志)
              snapshot  → 列出已生成的快照文件
EOF
  exit 1
}

case "$MODE" in
  dump|latest|shell|quorum|snapshot) ;;
  -h|--help) usage ;;
  *) echo "未知 mode: $MODE"; usage ;;
esac

echo ">>> 容器:$CONTAINER   模式:$MODE"
echo

# 检查容器是否在跑
if ! docker compose ps --format json 2>/dev/null | grep -q "$CONTAINER"; then
  if ! docker ps --format '{{.Names}}' | grep -q "^${CONTAINER}$"; then
    echo "❌ 容器 $CONTAINER 不在运行,先 docker compose up -d"
    exit 2
  fi
fi

run_in() {
  docker exec -i "$CONTAINER" bash -c "$1"
}

case "$MODE" in
  quorum)
    echo "=== Quorum Status ==="
    run_in "/opt/kafka/bin/kafka-metadata-quorum.sh \
              --bootstrap-server localhost:19092 describe --status"
    echo
    echo "=== Quorum Replication ==="
    run_in "/opt/kafka/bin/kafka-metadata-quorum.sh \
              --bootstrap-server localhost:19092 describe --replication"
    ;;

  snapshot)
    echo "=== 元数据目录与 snapshot 文件 ==="
    run_in "ls -lh $META_DIR/ | head -30"
    ;;

  shell)
    echo "=== 启动交互式 kafka-metadata-shell.sh ==="
    echo "(进去后可以执行 ls /  /  cat /topics/<name>/<partition>/data 等)"
    echo
    # 用最新一个 .log 段;如果有 .checkpoint(snapshot)也可以用
    LATEST_LOG=$(run_in "ls $META_DIR/*.log 2>/dev/null | sort | tail -1")
    if [[ -z "$LATEST_LOG" ]]; then
      echo "❌ 没找到 $META_DIR/*.log"; exit 3
    fi
    echo "使用文件: $LATEST_LOG"
    docker exec -it "$CONTAINER" /opt/kafka/bin/kafka-metadata-shell.sh \
      --snapshot "$LATEST_LOG"
    ;;

  latest|dump)
    if [[ "$MODE" == "latest" ]]; then
      FILES=$(run_in "ls $META_DIR/*.log 2>/dev/null | sort | tail -1")
    else
      FILES=$(run_in "ls $META_DIR/*.log 2>/dev/null | sort")
    fi
    if [[ -z "$FILES" ]]; then
      echo "❌ 在容器 $CONTAINER$META_DIR 下没找到 .log 文件"
      echo "   说明该节点不是 Controller Quorum 成员(普通 Broker 不存元数据日志副本)"
      echo "   Controller Quorum 一般是独立节点;如果你用了 combined 模式,"
      echo "   则 broker,controller 节点会有这个目录。"
      exit 4
    fi

    for f in $FILES; do
      echo "================================================================"
      echo "  解析: $f"
      echo "================================================================"
      run_in "/opt/kafka/bin/kafka-dump-log.sh \
                --cluster-metadata-decoder \
                --files $f \
                --print-data-log" \
        | head -300
      echo
    done

    cat <<EOF

读法说明:
  * baseOffset / lastOffset:本批次(RecordBatch)的 offset 区间;
  * payload 反序列化为 KRaft 定义的元数据事件类型;常见的有:
      RegisterBrokerRecord     —— Broker 加入集群
      UnregisterBrokerRecord   —— Broker 注销
      TopicRecord              —— 创建 Topic
      PartitionRecord          —— 创建分区时的初始状态
      PartitionChangeRecord    —— Leader/ISR/Replicas 变更
      ConfigRecord             —— 动态配置修改
      AccessControlEntryRecord —— ACL 变更
  * --cluster-metadata-decoder 是关键选项;不加它解出来都是十六进制乱码。

下一步:
  * 想交互查询:./dump_metadata_log.sh $CONTAINER shell
  * 想看 Quorum 健康:./dump_metadata_log.sh $CONTAINER quorum
EOF
    ;;
esac
python
#!/usr/bin/env python3
"""
inspect_metadata.py — 用 AdminClient 探测 Kafka 集群(KRaft 模式)的元数据

打印内容:
    1) cluster_id(整个集群的全局 UUID,KRaft 时代由 kafka-storage.sh 生成)
    2) Active Controller 节点 ID(即 KRaft Quorum 的 Raft Leader)
    3) 所有 Broker 列表(host:port + rack)
    4) 所有 Topic + 分区数 + 副本布局
    5) 默认配置的几个关键值(动态查询 broker config)

用法:
    pip install confluent-kafka
    python inspect_metadata.py
    python inspect_metadata.py --bootstrap localhost:9092
"""

from __future__ import annotations

import argparse
import sys

from confluent_kafka.admin import (
    AdminClient,
    ConfigResource,
    ResourceType,
)


def parse_args():
    p = argparse.ArgumentParser()
    p.add_argument("--bootstrap", default="localhost:9092",
                   help="bootstrap.servers")
    p.add_argument("--with-topics", action="store_true",
                   help="额外打印每个 Topic 的分区/副本布局")
    p.add_argument("--with-broker-configs", action="store_true",
                   help="额外打印每个 Broker 的 KRaft 关键配置")
    return p.parse_args()


def section(title):
    print()
    print("=" * 78)
    print(f"  {title}")
    print("=" * 78)


def main():
    args = parse_args()
    admin = AdminClient({"bootstrap.servers": args.bootstrap})

    md = admin.list_topics(timeout=15)

    section("CLUSTER")
    print(f"cluster_id          : {md.cluster_id}")
    print(f"controller broker_id: {md.controller_id}    "
          f"← 在 KRaft 集群里这是 Active Controller 在 Quorum 内的 node_id")
    print(f"alive broker count  : {len(md.brokers)}")

    section("BROKERS")
    print(f"{'id':<6}{'host':<30}{'port':<8}{'rack':<10}")
    print("-" * 60)
    for bid in sorted(md.brokers.keys()):
        b = md.brokers[bid]
        print(f"{bid:<6}{b.host:<30}{b.port:<8}{(b.rack or '-'):<10}")

    if args.with_topics:
        section("TOPICS")
        for tname in sorted(md.topics.keys()):
            t = md.topics[tname]
            if tname.startswith("__"):
                marker = "  [internal]"
            else:
                marker = ""
            print(f"\n{tname}{marker}")
            for pid in sorted(t.partitions.keys()):
                p = t.partitions[pid]
                print(f"    P{pid}  Leader={p.leader}  "
                      f"Replicas={list(p.replicas)}  ISR={list(p.isrs)}")

    if args.with_broker_configs:
        section("BROKER CONFIGS (KRaft 关键项)")
        keys_of_interest = [
            "process.roles",
            "node.id",
            "controller.quorum.voters",
            "controller.listener.names",
            "inter.broker.listener.name",
            "metadata.log.segment.bytes",
            "min.insync.replicas",
        ]
        for bid in sorted(md.brokers.keys()):
            res = ConfigResource(ResourceType.BROKER, str(bid))
            fut = admin.describe_configs([res])[res]
            try:
                conf = fut.result(timeout=10)
            except Exception as e:
                print(f"\n● Broker {bid}: 拉取配置失败:{e}")
                continue
            print(f"\n● Broker {bid}")
            for k in keys_of_interest:
                if k in conf:
                    v = conf[k]
                    src = v.source.name if hasattr(v.source, 'name') else v.source
                    print(f"    {k:<35} = {v.value}  ({src})")

    section("DONE")
    print("提示:")
    print("  * 在 KRaft 模式下,controller_id 是 Active Controller 的 node_id;")
    print("    要看完整 Quorum 状态用 kafka-metadata-quorum.sh describe --status。")
    print("  * Topic 名以 __ 开头是 Kafka 内部 Topic,包括:")
    print("      __consumer_offsets    (消费组 offset)")
    print("      __transaction_state   (事务协调)")
    print("      __cluster_metadata    (KRaft 元数据日志,仅 Quorum 节点能看到分区数据)")


if __name__ == "__main__":
    try:
        main()
    except KeyboardInterrupt:
        sys.exit(0)
txt
# =============================================================================
# kraft_config_example.properties
# -----------------------------------------------------------------------------
# Kafka KRaft 模式配置范例 —— 一份文件覆盖三种典型部署:
#
#   ① Combined 模式:单进程同时扮演 broker + controller   (学习 / 单机 / 小集群)
#   ② Controller-only:只做 Controller,不接客户端读写    (生产 Quorum 节点)
#   ③ Broker-only:只做 Broker,不参与 metadata Quorum    (生产业务节点)
#
# 用法:把下面对应模式的小节复制成 server.properties,按需修改 IP/端口/路径,
# 然后用 kafka-storage.sh format 初始化,再 kafka-server-start.sh 启动。
#
# 第一次启动前必须 format(KRaft 与 ZK 时代最大的区别之一):
#
#   # 在「集群里随便一台」生成 cluster ID
#   CLUSTER_ID=$(/opt/kafka/bin/kafka-storage.sh random-uuid)
#   echo $CLUSTER_ID
#
#   # 然后在「每个节点」上用同一个 ID format
#   /opt/kafka/bin/kafka-storage.sh format \
#       --config /opt/kafka/config/kraft/server.properties \
#       --cluster-id $CLUSTER_ID
#
# 启动:
#   /opt/kafka/bin/kafka-server-start.sh /opt/kafka/config/kraft/server.properties
# =============================================================================


# =============================================================================
# 模式 ① Combined(broker + controller 同进程)  —— 学习 / 小集群
# =============================================================================
process.roles=broker,controller

# 节点全局唯一 ID(KRaft 时代相当于 broker.id)
node.id=1

# Quorum 投票成员清单:写法 nodeId@host:port  (port 是 controller 端口)
# 三个节点的这一行配置必须**完全一致**,包括 host:port、id
controller.quorum.voters=1@kafka1:9093,2@kafka2:9093,3@kafka3:9093

# Quorum 内部使用哪个 listener 通信(必须出现在 listeners 里)
controller.listener.names=CONTROLLER

# 监听器:CONTROLLER 内部 RPC + PLAINTEXT 客户端入口
listeners=PLAINTEXT://0.0.0.0:9092,CONTROLLER://0.0.0.0:9093
advertised.listeners=PLAINTEXT://kafka1:9092
listener.security.protocol.map=CONTROLLER:PLAINTEXT,PLAINTEXT:PLAINTEXT
inter.broker.listener.name=PLAINTEXT

# 数据与元数据存放目录
log.dirs=/var/lib/kafka/data
metadata.log.dir=/var/lib/kafka/data

# 内部 Topic 副本数(与 Broker 数对齐,生产至少 3)
offsets.topic.replication.factor=3
transaction.state.log.replication.factor=3
transaction.state.log.min.isr=2
min.insync.replicas=2

# 教学环境关闭自动建 Topic,强制脚本/UI 显式创建
auto.create.topics.enable=false


# =============================================================================
# 模式 ② Controller-only (生产推荐:3 ~ 5 个独立 Controller 节点)
# -----------------------------------------------------------------------------
# 这种节点只跑 Raft Quorum,不接客户端读写,不存数据 Topic 副本。
# 把整个集群的「大脑」与「血管」物理隔离,避免业务流量影响选举与元数据同步。
# =============================================================================
# process.roles=controller
# node.id=3001
# controller.quorum.voters=3001@ctrl1:9093,3002@ctrl2:9093,3003@ctrl3:9093
# controller.listener.names=CONTROLLER
# listeners=CONTROLLER://0.0.0.0:9093
# # 注意:controller-only 节点不需要 advertised.listeners 给客户端
# # 也不需要 inter.broker.listener.name
# log.dirs=/var/lib/kafka/metadata
# metadata.log.dir=/var/lib/kafka/metadata
#
# # Controller 节点不参与业务数据,但仍然要为元数据日志预留磁盘
# # 推荐独立 SSD,至少 50GB
# metadata.log.segment.bytes=1073741824
# metadata.log.max.snapshot.interval.ms=3600000


# =============================================================================
# 模式 ③ Broker-only (生产推荐:N 个 Broker 节点,N 与业务规模相关)
# -----------------------------------------------------------------------------
# 这种节点专门跑数据 Topic 的 Leader/Follower,不参与 Quorum 投票。
# 它们通过 controller.quorum.voters 找到 Controller Quorum,注册自己、
# fetch __cluster_metadata 来获取元数据视图。
# =============================================================================
# process.roles=broker
# node.id=101
# controller.quorum.voters=3001@ctrl1:9093,3002@ctrl2:9093,3003@ctrl3:9093
# controller.listener.names=CONTROLLER
#
# listeners=PLAINTEXT://0.0.0.0:9092,INTERNAL://0.0.0.0:19092
# advertised.listeners=PLAINTEXT://broker-101.example.com:9092,INTERNAL://broker-101.internal:19092
# listener.security.protocol.map=CONTROLLER:PLAINTEXT,PLAINTEXT:PLAINTEXT,INTERNAL:PLAINTEXT
# inter.broker.listener.name=INTERNAL
#
# log.dirs=/data/kafka/log1,/data/kafka/log2     # 多盘并行
#
# # 业务数据相关
# num.partitions=3
# default.replication.factor=3
# min.insync.replicas=2
# unclean.leader.election.enable=false
# offsets.topic.replication.factor=3
# transaction.state.log.replication.factor=3
# transaction.state.log.min.isr=2
#
# # 副本相关
# num.replica.fetchers=4                  # 多线程拉副本
# replica.lag.time.max.ms=30000


# =============================================================================
# KRaft 通用调优(所有模式都可参考)
# =============================================================================
# 元数据快照:太频繁占 IO,太稀疏重启慢,1 小时 + 100MB 是一个平衡点
# metadata.log.max.snapshot.interval.ms=3600000
# metadata.log.segment.bytes=1073741824

# Quorum 选举超时:默认 1s。网络抖动较大的集群可调到 2~5s 减少误选举
# controller.quorum.election.timeout.ms=1000
# controller.quorum.fetch.timeout.ms=2000

# 元数据请求超时(普通 Broker 端使用)
# controller.socket.timeout.ms=30000

# 日志保留 / 清理
log.retention.hours=168
log.segment.bytes=1073741824
log.retention.check.interval.ms=300000

# 网络与线程
num.network.threads=3
num.io.threads=8
socket.send.buffer.bytes=102400
socket.receive.buffer.bytes=102400
socket.request.max.bytes=104857600

# JVM Heap:学习 1G,小集群 4G,大集群 6~8G
# 用 KAFKA_HEAP_OPTS="-Xms6G -Xmx6G" 在启动脚本里设置,不要写本文件


# =============================================================================
# 常见错误自检
# =============================================================================
# 1) 启动报 InconsistentClusterIdException:
#    每个节点的 meta.properties 里的 cluster.id 必须相同。
#    解决:rm -rf $log.dirs/meta.properties,重新 kafka-storage.sh format。
#
# 2) 启动报 InconsistentNodeIdException:
#    本节点 node.id 与 meta.properties 中记录的 node.id 不一致。
#    解决:要么修配置,要么删 meta.properties 重新 format。
#
# 3) 启动卡在「Awaiting socket connection to controller」:
#    controller.quorum.voters 配错;先 ping 通 host:port,再确认 voters 三节点完全一致。
#
# 4) Quorum 选不出 Leader:
#    多数派 = floor(N/2)+1。N=3 至少要 2 个节点活着。
#    用 kafka-metadata-quorum.sh describe --status 看当前 Voter 状态。
#
# 5) Active Controller 频繁切换:
#    可能 GC pause 太久(>1s 触发 election timeout),调大
#    controller.quorum.election.timeout.ms 或迁出业务流量。
# =============================================================================

dump_metadata_log.sh ↗ · inspect_metadata.py ↗ · kraft_config_example.properties ↗