【Kafka源码解读和使用指南】第54篇:Kafka控制器源码解析(四)——ZooKeeper Listener响应集群变化
上一篇【第53篇】Kafka控制器源码解析(三)——ReplicaStateMachine副本状态机
下一篇【第55篇】GroupCoordinator源码解析(一)——消费者组管理
摘要
Controller是Kafka集群的"大脑",但它不能主动"看到"集群里发生了什么——全靠ZooKeeper的Watcher机制来通知。当Broker上下线、Topic被创建删除、ISR发生变化时,ZK会触发对应的Listener,然后驱动状态机完成状态变更。
这一篇我们把Controller里所有ZK Listener都过一遍:BrokerChangeListener怎么感知Broker上下线、TopicChangeListener怎么处理Topic增删、IsrChangeNotificationListener怎么传播ISR变更、还有DeleteTopicsListener和PreferredReplicaElectionListener。读完这篇,你对"集群变化→Controller响应"的完整链路就彻底清楚了。
一、Listener体系总览——Controller的"感官系统"
Controller通过IZkChildListener和IZkDataListener两大接口监听ZK节点变化:
【Controller监听的ZK节点全景】
ZooKeeper 数据/子节点变化
┌──────────────────────────────────────────────────────┐
│ /brokers/ids (Broker上下线) │
│ │ │
│ └── BrokerChangeListener (IZkChildListener) │
│ │
│ /brokers/topics (Topic增删) │
│ │ │
│ └── TopicChangeListener (IZkChildListener) │
│ │
│ /brokers/topics/[topic]/partitions/[p]/state │
│ (分区Leader/ISR变化) │
│ │ │
│ └── PartitionModificationsListener (IZkDataListener)│
│ │
│ /isr_change_notification (ISR变更通知) │
│ │ │
│ └── IsrChangeNotificationListener (IZkChildListener)│
│ │
│ /admin/delete_topics (删除Topic请求) │
│ │ │
│ └── DeleteTopicsListener (IZkChildListener) │
│ │
│ /admin/preferred_replica_election │
│ (优先副本选举请求) │
│ │ │
│ └── PreferredReplicaElectionListener │
└──────────────────────────────────────────────────────┘
| Listener | 监听节点 | 接口类型 | 触发时机 |
|---|---|---|---|
| BrokerChangeListener | /brokers/ids |
IZkChildListener | Broker上下线 |
| TopicChangeListener | /brokers/topics |
IZkChildListener | Topic创建/删除 |
| PartitionModificationsListener | /brokers/topics/[t]/partitions/[p]/state |
IZkDataListener | 分区Leader/ISR数据变化 |
| IsrChangeNotificationListener | /isr_change_notification |
IZkChildListener | ISR变更通知 |
| DeleteTopicsListener | /admin/delete_topics |
IZkChildListener | 删除Topic请求 |
| PreferredReplicaElectionListener | /admin/preferred_replica_election |
IZkDataListener | 优先副本选举请求 |
二、BrokerChangeListener——感知Broker生死
这是最常用的Listener——每当你kafka-server-start.sh启动或kafka-server-stop.sh关闭一个Broker,这个Listener就会被触发。
// KafkaController.scala - BrokerChangeListener
class BrokerChangeListener extends IZkChildListener {
override def handleChildChange(parentPath: String,
currentBrokerIds: java.util.List[String]): Unit = {
inLock(controllerContext.controllerLock) {
if (hasStarted.get()) {
// ① 计算新增和死亡的Broker
val curBrokerIds = currentBrokerIds.map(_.toInt).toSet
val newBrokerIds = curBrokerIds -- controllerContext.liveBrokerIds
val deadBrokerIds = controllerContext.liveBrokerIds -- curBrokerIds
// ② 更新ControllerContext
newBrokerIds.foreach { id =>
val brokerInfo = ZkUtils.getBrokerInfo(zkUtils, id)
controllerContext.liveBrokers += brokerInfo
}
deadBrokerIds.foreach { id =>
controllerContext.liveBrokers -=
controllerContext.getBroker(id)
}
controllerContext.liveBrokerIds = curBrokerIds
// ③ 处理新Broker上线
if (newBrokerIds.nonEmpty) {
onBrokerStartup(newBrokerIds.toSeq)
}
// ④ 处理Broker下线
if (deadBrokerIds.nonEmpty) {
onBrokerFailure(deadBrokerIds.toSeq)
}
}
}
}
}
handleChildChange的核心逻辑分四步:
【Broker上线的处理链路】
Broker 3 启动
│
▼
/brokers/ids 子节点增加 "3"
│
▼
BrokerChangeListener.handleChildChange()
│
├── ① 更新ControllerContext.liveBrokers += Broker(3)
│
└── ② onBrokerStartup(Seq(3))
│
▼
⑴ replicaStateMachine.handleStateChanges()
→ 将Broker 3上的副本从Offline→OnlineReplica
│
▼
⑵ partitionStateMachine.triggerOnlinePartitionStateChange()
→ 将OfflinePartition→OnlinePartition(重新选举Leader)
│
▼
⑶ 向Broker 3发送UpdateMetadataRequest
→ 让它知道全集群的元数据
对于Broker下线,流程是反过来的:
// KafkaController.scala - onBrokerFailure()
def onBrokerFailure(deadBrokerIds: Seq[Int]): Unit = {
// ① 找出受影响的分区(Leader在这些死亡Broker上的)
val affectedPartitions = controllerContext.partitionsOnBrokers(deadBrokerIds)
// ② 将这些分区设为OfflinePartition → 触发Leader重新选举
partitionStateMachine.handleStateChanges(
affectedPartitions,
OfflinePartition,
offlinePartitionSelector)
// ③ 将死亡Broker上的副本设为OfflineReplica
deadBrokerIds.foreach { brokerId =>
val replicasOnBroker = controllerContext.replicasOnBroker(brokerId)
replicaStateMachine.handleStateChanges(
replicasOnBroker, OfflineReplica)
}
// ④ 更新ControllerContext
deadBrokerIds.foreach { id =>
controllerContext.liveBrokers -=
controllerContext.getBroker(id)
}
}
三、TopicChangeListener——Topic增删的"守门员"
当你执行kafka-topics.sh --create或--delete时,这个Listener被触发。
// KafkaController.scala - TopicChangeListener
class TopicChangeListener extends IZkChildListener {
override def handleChildChange(parentPath: String,
currentTopics: java.util.List[String]): Unit = {
inLock(controllerContext.controllerLock) {
if (hasStarted.get()) {
val curTopics = currentTopics.toSet
// ① 计算新增和删除的Topic
val newTopics = curTopics -- controllerContext.allTopics
val deletedTopics = controllerContext.allTopics -- curTopics
// ② 更新ControllerContext
controllerContext.allTopics = curTopics
// ③ 处理新增Topic
if (newTopics.nonEmpty) {
val addedPartitionReplicaAssignment =
ZkUtils.getReplicaAssignmentForTopics(zkUtils, newTopics.toSeq)
controllerContext.partitionReplicaAssignment ++=
addedPartitionReplicaAssignment
onNewTopicCreation(newTopics,
addedPartitionReplicaAssignment.keySet)
}
// ④ 处理删除的Topic
if (deletedTopics.nonEmpty) {
onTopicDeletion(deletedTopics)
}
}
}
}
}
新建Topic的完整处理流程:
【创建Topic的处理链路】
kafka-topics.sh --create --topic test
│
▼
ZK /brokers/topics 下新增子节点 "test"
│
▼
TopicChangeListener.handleChildChange()
│
├── ① 从ZK加载AR集合
│ ZkUtils.getReplicaAssignmentForTopics()
│
└── ② onNewTopicCreation()
│
▼
⑴ replicaStateMachine.handleStateChanges()
→ 将所有副本从NonExistentReplica→NewReplica
│
▼
⑵ partitionStateMachine.handleStateChanges()
→ 将所有分区从NonExistentPartition→NewPartition→OnlinePartition
→ 选举初始Leader
│
▼
⑶ 向所有Broker发送UpdateMetadataRequest
四、PartitionModificationsListener——监听分区状态变化
每个分区在ZK上都有一个/brokers/topics/[topic]/partitions/[partition]/state节点,存储着当前的Leader和ISR信息。当ISRF发生变化时,这个节点的数据会被更新,此Listener被触发。
// KafkaController.scala - PartitionModificationsListener
class PartitionModificationsListener(topic: String, partitionId: Int)
extends IZkDataListener {
override def handleDataChange(dataPath: String,
data: Object): Unit = {
inLock(controllerContext.controllerLock) {
// ① 从ZK读取最新的LeaderAndIsr信息
val newLeaderAndIsr = ZkUtils.getLeaderAndIsrForPartition(
zkUtils, topic, partitionId)
// ② 更新ControllerContext
val leaderIsrAndControllerEpoch =
LeaderIsrAndControllerEpoch(newLeaderAndIsr, controller.epoch)
controllerContext.partitionLeadershipInfo.put(
TopicAndPartition(topic, partitionId),
leaderIsrAndControllerEpoch)
// ③ 如果Leader发生了变化,发送UpdateMetadataRequest
// (实际上ReplicaManager收到LeaderAndIsrRequest后会自己处理)
}
}
}
五、IsrChangeNotificationListener——ISR变更广播
当ISR发生变化时,ReplicaManager会将变更记录到/isr_change_notification节点下,所有Broker都会收到通知。
【ISR变更通知流程】
Broker 1 (Leader of P0)
ISR从 [1,2,3] 变为 [1,2] (Broker 3滞后)
│
▼
ReplicaManager 写入 /isr_change_notification/isr_change_00001
│
▼
所有Broker上的 IsrChangeNotificationListener 被触发
│
▼
调用 replicaStateMachine.handleStateChanges()
重新加载ISR信息
六、DeleteTopicsListener——Topic删除的"执行者"
当你执行kafka-topics.sh --delete时,删除请求被写入/admin/delete_topics,此Listener被触发。
// KafkaController.scala - DeleteTopicsListener
class DeleteTopicsListener extends IZkChildListener {
override def handleChildChange(parentPath: String,
children: java.util.List[String]): Unit = {
inLock(controllerContext.controllerLock) {
val topicsToBeDeleted = children.toSet
val topicsWithDeletionStarted =
controllerContext.topicsToBeDeleted.filter { topic =>
controller.deleteTopicManager.isTopicDeletionInProgress(topic)
}
// ① 找出新标记的待删除Topic
val newTopicsToBeDeleted =
topicsToBeDeleted -- topicsWithDeletionStarted
// ② 启动删除流程
controller.deleteTopicManager.enqueueTopicsForDeletion(
newTopicsToBeDeleted.toSeq)
}
}
}
Topic删除是一个多步骤的流程:
【Topic删除的完整流程】
kafka-topics.sh --delete --topic test
│
▼
/admin/delete_topics 写入 "test"
│
▼
DeleteTopicsListener
│
▼
deleteTopicManager.enqueueTopicsForDeletion()
│
├── ① 将所有分区设为OfflinePartition
├── ② 向所有Broker发送StopReplicaRequest(delete=true)
│ → 各Broker删除本地日志文件
├── ③ 从ZK删除 /brokers/topics/test
└── ④ 更新ControllerContext,移除该Topic
七、各Listener的协同关系
所有这些Listener不是孤立工作的,它们共同构成了Controller响应集群变化的"神经系统":
【集群变化 → Controller响应的完整链路】
事件 Listener 状态机操作
───────────────────────────────────────────────────────────────────
Broker上线 → BrokerChangeListener → 副本Online + 分区Online
Broker下线 → BrokerChangeListener → 副本Offline + Leader重选
Topic创建 → TopicChangeListener → 副本New + 分区Online
Topic删除 → DeleteTopicsListener → 副本Deletion + 分区NonExistent
ISR变化 → IsrChangeNotification → 更新LeaderAndIsr缓存
Listener
优先副本选举 → PreferredReplica → 分区OnlinePartition(重选Leader)
ElectionListener
本篇小结
- ZK监听是Controller感知集群变化的唯一途径:所有集群状态变化都通过ZK Watcher通知Controller
- 六大Listener各司其职:BrokerChangeListener管Broker生死、TopicChangeListener管Topic增删、PartitionModificationsListener管分区状态、IsrChangeNotificationListener管ISR变更广播、DeleteTopicsListener管Topic删除、PreferredReplicaElectionListener管优先副本选举
- Listener驱动状态机:每个Listener被触发后,最终都调用PartitionStateMachine或ReplicaStateMachine来完成状态转换
- handleChildChange vs handleDataChange:前者监听子节点增删(Broker、Topic),后者监听数据内容变化(分区状态)
下一篇我们将走进GroupCoordinator——消费者组的"大管家",看它如何管理消费者组的生命周期和offset提交。
上一篇【第53篇】Kafka控制器源码解析(三)——ReplicaStateMachine副本状态机
下一篇【第55篇】GroupCoordinator源码解析(一)——消费者组管理
更多推荐


所有评论(0)