这是《Kafka 零基础实战指南》的独立章节版。本章从概念、实操和生产排查三个视角展开,代码块保留了原书可直接运行的版本。 Kafka 集群里有两类“管理者”:Controller 管集群元数据与分区 Leader,Group Coordinator 管消费组成员与位移。理解它们,rebalance、选主、位移提交这些行为就不再神秘。
12.1 Controller 是谁,做什么
Controller 是集群元数据的仲裁者,负责:
- Broker 上下线感知;
- 分区 Leader 选举(第 11 章);
- topic 的创建、删除、配置变更;
- 分区副本重分配(reassignment);
- 向所有 Broker 广播元数据更新。
在 ZooKeeper 时代,Controller 通过抢占临时节点选举,元数据存 ZooKeeper;KRaft 时代,Controller 是一个 Raft 仲裁组,元数据就是 __cluster_metadata 这条日志本身。
12.2 KRaft 下的 Controller
KRaft 集群中常见两种部署形态:
combined(合并模式):
broker1(broker+controller) broker2(broker+controller) broker3(broker+controller)
适合中小集群,省机器
separated(分离模式):
controller1 controller2 controller3 <- 只跑仲裁
broker1 ... brokerN <- 只存数据
适合大规模或元数据敏感场景
由 process.roles 决定:
process.roles=broker,controller # combined
# 或
process.roles=broker # 分离模式的数据节点
Controller 数量必须是奇数(3 或 5),因为 Raft 需要多数派(quorum)达成共识。元数据变更流程:
- 客户端/Broker 发起变更(建 topic、Broker 注册);
- Leader Controller 把变更追加到元数据日志;
- Raft 多数派确认后 commit;
- 所有 Broker(作为 learner)拉取并应用新元数据。
查看仲裁状态:
bin/kafka-metadata-quorum.sh \
--bootstrap-server localhost:9092 describe --status
输出包含 LeaderId、LeaderEpoch、HighWatermark,可以直观看到 Raft 进度。
12.3 Group Coordinator 的职责
每个消费组在某个 Broker 上有一个 Group Coordinator,负责:
- 处理
JoinGroup/SyncGroup(加入组、接收分配方案); - 维护成员与心跳(判断谁掉线);
- 触发与管理 Rebalance;
- 接收并存储位移提交(写入
__consumer_offsets)。
哪个 Broker 协调哪个组?由组名的哈希决定:
coordinator 分区 = abs(group.hashCode()) % 50
该分区的 Leader 所在 Broker = 协调器
(50 是 offsets.topic.num.partitions 默认值。)所以组名固定,协调器位置也固定;该分区 Leader 切换时,协调器随之切换,客户端会自动重新发现。
12.4 消费组协议:一次完整的组生命周期
1. FindCoordinator 消费者问:我的协调器在哪台 Broker?
2. JoinGroup 所有成员向协调器报名,选出一个 leader 消费者
3. SyncGroup 组 leader 把分配方案发给协调器,协调器下发给成员
4. Heartbeat 成员定期心跳维持成员资格
5. OffsetCommit/OffsetFetch 提交与拉取位移
6. LeaveGroup 成员主动退出(优雅停机)
注意第 2、3 步中的“组 leader”是消费者 leader(不是 Broker):它运行分配策略,决定哪个成员拿哪些分区。协调器只负责组织流程与存储结果。这就是为什么切换分配策略只需要改客户端配置(第 15 章)。
12.5 __consumer_offsets 内部结构
位移提交写入内部 compact topic __consumer_offsets,key 是 (group, topic, partition):
key: group.id + topic + partition
value: offset + metadata + timestamp
配合 compact 清理,每个 key 只保留最新提交,因此该主题体积可控。读取它需要专用 formatter(第 4 章命令)。两点实践提醒:
- 生产集群该主题 RF 必须是 3(默认配置项
offsets.topic.replication.factor),否则协调器所在 Broker 宕机会影响大量消费组; - 重置位移命令本质上也是往这个主题写新记录,这就是为什么执行前要停止消费者。
12.6 事务协调器
事务生产者初始化时,会根据 transactional.id 定位到某个 Broker 上的 TransactionCoordinator:
- 事务状态持久化在
__transaction_state; - 协调器负责分配 producer epoch、两阶段提交(PrepareCommit -> 写 marker -> CompleteCommit);
- Broker 端用 LSO(Last Stable Offset)控制
read_committed消费者能读到哪(第 14 章)。
12.7 元数据如何到达客户端
生产/消费客户端会缓存集群元数据(分区 Leader 位置)。刷新时机:
- 定期:
metadata.max.age.ms(默认 5 分钟); - 被动:请求收到
NOT_LEADER_OR_PARTITION、UNKNOWN_TOPIC_OR_PARTITION等错误时立刻刷新。
所以 Leader 切换后的短暂窗口内,客户端可能仍把请求发给旧 Leader,收到错误后自动纠正。这就是“切主瞬间偶发报错又自愈”的原因。
本章小结
- Controller 管元数据与分区 Leader;KRaft 用 Raft 日志
__cluster_metadata承载元数据; - Group Coordinator 管消费组成员、rebalance 与位移,位置由 group 哈希定位;
- 消费组里的“组 leader”(一个消费者)负责计算分区分配方案;
- 位移存储在 compact 的
__consumer_offsets,事务状态在__transaction_state; - 客户端元数据缓存 + 错误触发刷新,让 Leader 切换对业务基本透明。
思考题
- 为什么 KRaft Controller 数量要配成奇数?2 个 Controller 会发生什么?
__consumer_offsets为什么用 compact 而不是普通 delete?- 消费组 rebalance 过程中,协调器和“组 leader 消费者”各做什么?