上一篇【第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源码解析(一)——消费者组管理


Logo

汇聚全球AI编程工具,助力开发者即刻编程。

更多推荐