ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

Kafka 分区分配与 Leader 选举:Controller 如何用 ZK/KRaft 元数据驱动状态机

Kafka 分区分配与 Leader 选举:Controller 如何用 ZK/KRaft 元数据驱动状态机 Kafka 分区分配与 Leader 选举Controller 如何用 ZK/KRaft 元数据驱动状态机1. 先看一个真实困惑为什么副本都在分区却不可用假设你负责一个订单履约系统Kafka 集群有 6 个 Broker某个核心 Topic 设置 3 副本。某天凌晨一台 Broker 宕机监控告警显示“分区不可用”但你登上去看一眼发现 3 个副本所在的机器里还有两台活着数据盘也正常。为什么分区还是不能读写这类问题的答案通常不在“副本够不够”而在“谁来决定哪个副本当 Leader、这个决定有没有成功写进元数据”。Kafka 的可用性不是 Broker 之间自发协商出来的而是由一个特殊角色——集群里唯一负责管理元数据的 Broker也就是 Controller——把“哪个分区归谁、谁是 Leader、哪些副本在同步”这些信息写进一份权威元数据里其他 Broker 再据此行动。如果你只记住“Kafka 有副本所以高可用”就会在 Leader 选举卡住、Controller 频繁切换、ISR 反复抖动时无从下手。本文要建立的正是这条链路元数据从哪里来 → Controller 怎么当选 → 分区状态怎么流转 → Leader 怎么被选出 → 客户端怎么拿到并缓存这些信息。先记住一个最小模型Controller 是写元数据的人分区状态机是元数据的状态定义Leader 选举是状态机的一次跃迁元数据缓存是所有人读到最新状态的途径。ZK 模式和 KRaft 模式的差别本质上只是“这份元数据存在哪、谁来写”。2. 整体框架角色、元数据与一次决策的流转把整体拆成四部分Broker 集合、Controller、元数据存储、客户端。Broker 是干活的人负责收发消息和复制数据Controller 是管账的人负责分配分区、选 Leader、推进状态元数据存储是账本ZK 模式下是 ZooKeeper 的 znodeKRaft 模式下是内部的元数据日志客户端是办事的群众通过元数据缓存知道该找谁。一次“分区状态变化”的完整流转顺序是这样的某个触发事件发生创建 Topic、Broker 上下线、副本落后→ Controller 感知事件 → 计算新的分配或 Leader → 把新状态写入元数据存储 → 通知相关 Broker → Broker 加载本地副本状态 → 客户端刷新元数据缓存 → 请求被路由到新 Leader。触发事件建 Topic / Broker 下线 / ISR 收缩 | v --------------------- 写状态 ----------------------- | Controller | --------------- | 元数据存储 | | 分区状态机 / 选举逻辑 | --------------- | ZK znode 或 KRaft 日志 | --------------------- 读状态 ----------------------- | | | LeaderAndIsr / UpdateMetadata | 元数据变更事件 v v --------------------- ----------------------- | 目标 Broker 副本 | | 所有 Broker 缓存 | --------------------- ----------------------- | | v v 成为新 Leader / Follower 客户端元数据缓存刷新这里最容易误解的是Controller 并不搬运消息也不参与数据复制。它只做“管理面”的决策。真正把数据从 Leader 复制到 Follower 的是副本拉取机制。把管理面和数据面分开是理解后续所有状态变化的前提。3. Controller 选举谁有资格管账怎么避免两个管账人Controller 的核心约束是“同一时刻只有一个”。如果出现两个 Controller 各自写元数据就会出现“脑裂”一台认为分区 A 的 Leader 是 Broker1另一台认为是 Broker2客户端被指向两个不同的 Leader数据就乱了。所以 Controller 选举的第一目标是互斥第二目标才是快速失败转移。在 ZK 模式下Controller 选举依赖 ZooKeeper 的临时节点和顺序性。每个 Broker 启动时都会尝试在/controller路径创建临时节点谁先创建成功谁就是 Controller其余 Broker 在该节点上注册监听。当 Controller 所在 Broker 宕机或会话超时临时节点消失所有 Broker 收到通知重新竞争创建。临时节点的“会话绑定”特性天然保证了互斥会话没了节点就没了不会留下幽灵 Controller。在 KRaft 模式下没有 ZooKeeperController 由一组专门的 Controller 节点通常是 3 个或 5 个组成通过 Raft 协议在内部选举出一个 Active Controller其余是 Standby。Raft 的多数派投票和任期号term保证了同一任期只有一个 Leader且日志按顺序复制。这个变化的意义是把“元数据一致性”从外部系统收回到 Kafka 自己手里减少了运维组件也让元数据变更的延迟更可控。对比项ZK 模式KRaft 模式元数据存储ZooKeeper znode内部__cluster_metadata日志Controller 数量1 个由 Broker 竞争产生一组 Controller 节点选出 1 个 Active选举机制临时节点抢占 Watch 通知Raft 任期投票 日志复制故障转移速度受 ZK 会话超时影响通常秒级到十几秒通常更快取决于 Controller 心跳配置运维组件需要独立维护 ZK 集群只需 Kafka 自身元数据规模瓶颈分区数大时 znode 压力明显日志追加扩展性更好设计取舍很清楚ZK 模式成熟、生态工具多但多了一套要维护的系统且元数据规模受 ZK 写放大限制KRaft 模式简化了部署、提升了元数据扩展性但要求 Controller 节点独立规划资源且老版本工具链兼容性需要确认。当集群规模小、团队已有成熟 ZK 运维时ZK 模式仍然可用当分区数上万、希望减少组件时KRaft 是更合理的方向。4. 元数据到底存了什么分区状态机与副本状态机Controller 管账时账本上的每一行不是随便写的而是有明确状态定义的。分区状态机描述一个分区整体处于哪个阶段副本状态机描述某个副本在这个分区里扮演什么角色。两者是“整体”和“局部”的关系。分区状态机常见的状态包括NonExistent分区还不存在、NewPartition刚创建、还没分配 Leader、OnlinePartition有可用 Leader能读写、OfflinePartition没有可用 Leader不能读写。关键跃迁是NewPartition → OnlinePartition选出一个 Leader和OnlinePartition → OfflinePartitionLeader 挂掉且没有合格副本可接。副本状态机常见状态包括NewReplica新副本正在同步、OnlineReplica正常在线、OfflineReplica所在 Broker 下线、Leader当前是主副本、Follower跟随 Leader 复制。一个分区可能有多个OnlineReplica但同一时刻最多一个Leader。分区状态含义是否可读写典型进入条件NonExistent分区尚未创建否删除后或创建前NewPartition已创建未选主否收到 CreateTopics 请求OnlinePartition有有效 Leader是选出 Leader 并通知成功OfflinePartition无有效 Leader否Leader 下线且无合格副本这里有一个容易忽略的点状态机不是“自动机自己跑”而是 Controller 在特定事件下主动调用状态转移函数。比如OnlinePartition → OfflinePartition并不是定时扫描发现的而是 Controller 在处理 Broker 下线事件时判断“这个分区的所有 ISR 副本都不可用”才执行转移。所以状态变化的实时性取决于事件感知速度。5. Leader 选举从 ISR 里挑一个还是等 ISR 恢复Leader 选举的核心规则可以用一句话概括从 ISRIn-Sync Replicas保持同步的副本集合里选一个存活副本当 Leader。ISR 是“数据不落后太多”的副本名单只有名单里的副本才有资格当 Leader因为它们的数据最接近已提交状态。选举的具体过程是Controller 得知某分区 Leader 不可用 → 检查该分区的 ISR → 如果 ISR 中有存活副本选其中一个通常按 AR 或特定顺序优先成为新 Leader其余 ISR 副本成为 Follower → 把新的 LeaderAndIsr 信息写入元数据并通知相关 Broker → 分区回到OnlinePartition。如果 ISR 中没有任何存活副本分区进入OfflinePartition此时如果开启unclean.leader.election.enable才可能从非 ISR 副本中选主。Leader 不可用 | v 检查该分区 ISR | ---- ISR 中有存活副本 ---- 是 ---- 选一个为 Leader | | | v | 写入 LeaderAndIsr | | | v | 分区 OnlinePartition | ---- 否 ---- unclean.leader.election 是否开启 | -- 是 -- 从非 ISR 选主可能丢数据 | -- 否 -- 分区 OfflinePartition等待 ISR 恢复unclean.leader.election.enable是一个典型的取舍开启它分区能在 ISR 全挂时尽快恢复可用但新 Leader 可能缺少已提交消息造成数据丢失关闭它默认且更安全分区会保持不可用直到合格副本回来。对于订单、支付类场景宁可短时不可用也不要丢数据对于日志采集类场景可用性优先时可以考虑开启但必须接受数据回退。6. 一次创建 Topic 的完整推演现在让一次创建 Topic 的请求完整走一遍看看 Controller 在每一步做了什么。输入是客户端发出CreateTopicsRequest(topicorders, partitions3, replicationFactor3)。第一步任意 Broker 收到请求如果是 KRaft 模式会把请求转发给 Active Controller如果是 ZK 模式则转发给当前 Controller。第二步Controller 校验 Topic 是否已存在、副本因子是否合法、Broker 数量是否够用。第三步Controller 执行分区分配算法默认是轮询加机架感知为 3 个分区各挑 3 个副本并确定每个分区的首选 Leader。第四步Controller 把 Topic 和分区的初始状态写入元数据存储分区先是NewPartition。第五步Controller 向涉及的 Broker 发送LeaderAndIsr请求Broker 创建本地副本目录并进入副本状态机的相应状态。第六步选出的 Leader 副本开始接受生产请求分区进入OnlinePartition。importorg.apache.kafka.clients.admin.Admin;importorg.apache.kafka.clients.admin.AdminClientConfig;importorg.apache.kafka.clients.admin.NewTopic;importorg.apache.kafka.common.errors.TopicExistsException;importjava.util.Collections;importjava.util.Properties;importjava.util.concurrent.ExecutionException;publicclassCreateTopicDemo{publicstaticvoidmain(String[]args){PropertiespropsnewProperties();props.put(AdminClientConfig.BOOTSTRAP_SERVERS_CONFIG,localhost:9092);props.put(AdminClientConfig.REQUEST_TIMEOUT_MS_CONFIG,5000);try(AdminadminAdmin.create(props)){NewTopictopicnewNewTopic(orders-demo,3,(short)3);try{admin.createTopics(Collections.singleton(topic)).all().get();System.out.println(topic created: orders-demo);}catch(ExecutionExceptione){if(e.getCause()instanceofTopicExistsException){System.out.println(topic already exists, skip);}else{throwe;}}}catch(Exceptione){System.err.println(create topic failed: e.getMessage());}}}这个示例的目标是观察创建 Topic 后元数据的实际分布。前置环境是一个可用的 Kafka 集群单机 3 个 Broker 或 1 个 Broker 均可1 个 Broker 时副本因子要降到 1。运行后预期输出topic created: orders-demo。关键点是Admin实现了AutoCloseable用 try-with-resources 保证连接释放all().get()会把异步请求阻塞到完成方便演示生产代码里通常不这样阻塞。容易改错的地方是副本因子大于 Broker 数会报InvalidReplicationFactorException。7. 一次 Leader 宕机的完整推演第二个场景更贴近线上orders-demo的 partition 0 的 Leader 在 Broker1 上ISR 是[1,2,3]。Broker1 进程崩溃。Controller 通过 Broker 会话超时感知到 Broker1 下线进入故障处理流程。第一步Controller 把 Broker1 上所有副本标记为OfflineReplica。第二步对 partition 0检查 ISR[1,2,3]Broker1 已下线存活的是 2 和 3。第三步Controller 从存活 ISR 里选 Broker2 为新 Leader写入LeaderAndIsr(partition0, leader2, isr[2,3])。第四步通知 Broker2 和 Broker3Broker2 从 Follower 变为 LeaderBroker3 继续跟随。第五步客户端在元数据刷新周期后知道新 Leader 是 Broker2生产请求重新路由。时刻 T0: partition0 LeaderBroker1, ISR[1,2,3] 时刻 T1: Broker1 崩溃会话超时 时刻 T2: Controller 标记 Broker1 副本 OfflineReplica 时刻 T3: ISR 收缩为 [2,3]选 Broker2 为 Leader 时刻 T4: 写入 LeaderAndIsr通知 Broker2/Broker3 时刻 T5: 客户端元数据刷新请求转向 Broker2这个推演解释了第 1 节里的困惑如果存活副本不在 ISR 里比如 ISR 只有[1]Broker1 一挂ISR 里没有存活副本分区就直接进入OfflinePartition。副本“物理上活着”不等于“有资格当 Leader”资格由 ISR 决定。ISR 收缩的触发条件由replica.lag.time.max.ms默认 30 秒控制副本超过这个时间没追上 Leader就会被移出 ISR这直接影响故障时有多少候选 Leader。8. 元数据缓存客户端和 Broker 各自记什么元数据如果每次请求都去问 ControllerController 会被压垮延迟也不可接受。所以 Kafka 在两端都做了缓存Broker 缓存集群元数据用于处理请求客户端缓存 Topic 的 Leader 信息用于路由。客户端的元数据缓存是“按需刷新”的生产者第一次发送某 Topic 消息时如果本地没有该 Topic 的元数据就向任意 Broker 发MetadataRequest获取之后按metadata.max.age.ms默认 5 分钟定期刷新。当请求遇到NotLeaderForPartitionException或UnknownTopicOrPartitionException时客户端会立即触发一次元数据刷新并重试。这解释了为什么 Leader 切换后客户端通常能在很短时间内恢复而不是等满 5 分钟。Broker 端的元数据缓存则由 Controller 通过UpdateMetadataRequest主动推送。Controller 每次状态变化后都会把最新的元数据快照发给所有 BrokerBroker 更新本地缓存。这里有一个边界如果某个 Broker 长期收不到更新比如网络分区它的缓存会变旧可能把请求转发给已经不在 ISR 的旧 Leader客户端据此重试最终靠刷新纠正。缓存位置缓存内容更新方式过期/异常处理客户端Topic 分区 → Leader/副本按需拉取 定期刷新遇到 NotLeader 立即刷新重试Broker全量集群元数据Controller 主动推送收不到更新则缓存变旧Controller权威元数据自己写入依赖元数据存储持久化9. ZK 与 KRaft同一套状态机两种元数据载体理解 ZK 和 KRaft 的差异不要从“哪个新”出发而要从“元数据写在哪、读路径多长”出发。状态机、选举规则、ISR 逻辑在两种模式下基本一致变化的是元数据的持久化和传播层。在 ZK 模式下Controller 把状态写入 znode其他 Broker 通过 Watch 感知变化。Watch 是一次性触发的触发后要重新注册这就带来一个固有延迟状态变化传播要经过“写 znode → 触发 Watch → Broker 重新读”的多跳。分区数很大时znode 数量和 Watch 风暴会让 Controller 切换变慢这是 ZK 模式的主要痛点。在 KRaft 模式下元数据是一份追加写的日志存在__cluster_metadataTopic 里。Active Controller 把元数据变更作为记录追加到日志Standby Controller 复制日志Broker 通过拉取或推送获得元数据快照。由于是顺序追加写放大比 znode 小元数据规模扩展性好故障转移也不再依赖 ZK 会话超时。代价是 Controller 节点需要独立规划且元数据日志的磁盘 I/O 要稳定。# 查看 KRaft 模式下的元数据日志 Topic需要能访问内部 Topickafka-topics.sh --bootstrap-server localhost:9092\--describe--topic__cluster_metadata# 查看某分区的 Leader 与 ISR 状态kafka-topics.sh --bootstrap-server localhost:9092\--describe--topicorders-demo# 查看 Controller 当前是谁kafka-metadata-quorum.sh --bootstrap-server localhost:9092 describe--status这三条命令的作用分别是确认元数据日志存在、查看分区 Leader 与 ISR、确认 Active Controller。kafka-metadata-quorum.sh是 KRaft 专有工具ZK 模式下要用zookeeper-shell.sh查看/controller和/brokers路径。容易改错的地方是端口和引导地址--bootstrap-server需要指向任意可用 Broker而不是 Controller 专用端口除非显式配置。10. 用管理客户端观测选举与状态变化第三个完整示例是一个可复现的观测脚本它周期性打印 Topic 的分区 Leader 和 ISR用来在故障演练中观察状态机跃迁。前置环境是 Kafka 集群和一个已存在的 Topicorders-demo。importorg.apache.kafka.clients.admin.Admin;importorg.apache.kafka.clients.admin.AdminClientConfig;importorg.apache.kafka.clients.admin.TopicDescription;importorg.apache.kafka.common.TopicPartitionInfo;importjava.util.Collections;importjava.util.Properties;publicclassWatchLeaderAndIsr{publicstaticvoidmain(String[]args)throwsException{Stringtopicargs.length0?args[0]:orders-demo;PropertiespropsnewProperties();props.put(AdminClientConfig.BOOTSTRAP_SERVERS_CONFIG,localhost:9092);try(AdminadminAdmin.create(props)){for(inti0;i10;i){TopicDescriptiondescadmin.describeTopics(Collections.singleton(topic)).allTopicNames().get().get(topic);for(TopicPartitionInfop:desc.partitions()){System.out.printf(partition%d leader%d isr%s replicas%s%n,p.partition(),p.leader().id(),p.isr(),p.replicas());}Thread.sleep(2000);}}}}运行方式是先编译再传入 Topic 名执行。预期输出形如partition0 leader2 isr[2, 3] replicas[1, 2, 3]。关键观察点是当 Leader 所在 Broker 被停掉后leader字段会在若干秒后变成另一个 Brokerisr会收缩重启该 Broker 后它会先从 Follower 开始同步ISR 再扩张。这个脚本适合在预发环境做故障演练不适合生产高频调用因为describeTopics是管理面请求频繁调用会给 Controller 增加压力。11. 常见误区看起来相同机制上完全不同第一个误区是“副本在线就等于分区可用”。副本所在进程活着只能说明它可能参与复制能不能当 Leader取决于它是否在 ISR。如果它落后太多被踢出 ISR即使进程健康也没有资格。第二个误区是“Controller 挂了集群就完了”。Controller 挂了只影响管理面已经达到OnlinePartition的分区仍可正常读写直到需要新的管理决策比如再挂一个 Broker。真正危险的是 Controller 长时间无法选出导致故障无法修复。第三个误区是“ISR 越大越安全”。ISR 大意味着候选 Leader 多但每条消息都要等更多副本确认取决于acks吞吐会下降。ISR 大小是可用性和延迟的平衡不是越大越好。第四个误区是“KRaft 只是把 ZK 换掉行为一样”。元数据载体变了故障转移速度、元数据扩展性、运维工具都变了。升级前要确认版本、工具链和 Controller 节点规划不能假设行为完全等价。第五个误区是“客户端元数据缓存 5 分钟才刷新Leader 切换要等 5 分钟”。实际上遇到NotLeaderForPartition会立即刷新5 分钟只是兜底的定期刷新周期。12. 生产实践建议把机制变成配置和流程第一Controller 节点要独立规划。ZK 模式下 Controller 和 Broker 混部时要保证 Controller 所在 Broker 不被数据流量打满KRaft 模式下建议用独立 Controller 节点并给元数据日志盘留足 IOPS。第二谨慎调整replica.lag.time.max.ms。调大它ISR 更稳定但故障时可能选到较落后的副本调小它ISR 收缩更敏感可能频繁抖动。默认 30 秒对多数场景够用延迟敏感系统可以适当调小但要配合监控。第三unclean.leader.election.enable默认关闭。只有明确接受数据丢失的日志类场景才考虑开启且要单独评估。第四副本因子至少 3并让副本分散在不同机架配置broker.rack。这直接决定 ISR 在机架故障时是否还有存活副本。第五监控要覆盖 Controller 切换次数、ISR 收缩扩张频率、离线分区数、UnderReplicatedPartitions。这些指标比单纯的 Broker CPU 更能反映管理面健康度。第六升级 KRaft 要按官方迁移路径走先做元数据快照和演练不要直接在核心集群上试。13. 排障清单分区不可用时按这个顺序查当收到“分区不可用”告警时建议按以下顺序排查每一步都有明确的判断依据。第一步确认当前 Controller 是否存在且稳定。ZK 模式看/controller节点是否存在、是否有频繁切换KRaft 模式用kafka-metadata-quorum.sh describe --status看 Active Controller 和任期。如果 Controller 不稳定先解决它其余状态都不准。第二步查目标分区的 Leader 和 ISR。用kafka-topics.sh --describe看Leader是否为 -1、ISR 是否为空或只剩已下线副本。Leader: -1基本等于分区处于OfflinePartition。第三步确认 ISR 中副本所在 Broker 是否存活。如果 ISR 里的 Broker 全挂了分区不可用是设计使然不是 bug。第四步检查replica.lag.time.max.ms和副本同步是否被拖慢。磁盘慢、网络抖动、Broker 负载高都会让副本追不上导致被踢出 ISR。第五步检查 Controller 日志中是否有LeaderAndIsr发送失败或状态转移异常这些日志通常直接指出卡在哪一步。第六步如果刚做过扩缩容确认分区分配是否完成、新副本是否已追上并加入 ISR。14. 面试/复盘问题检验你是否真的理解Controller 选举为什么必须保证互斥如果出现两个 Controller客户端会看到什么现象分区状态机从NewPartition到OnlinePartition需要经历哪些步骤哪一步失败会导致分区一直不可用一个副本进程健康但不在 ISR它能被选为 Leader 吗为什么unclean.leader.election.enable开启后恢复可用性和数据一致性之间的具体代价是什么客户端元数据缓存多久刷新一次遇到NotLeaderForPartition时会发生什么ZK 模式下元数据传播要经过哪些环节为什么分区数多时 Controller 切换会变慢KRaft 模式下元数据日志和普通消息日志在机制上有什么相同与不同如果 ISR 频繁抖动你会先看哪些指标、调哪些配置15. 总结把分散机制收回一张决策框架回到最初的困惑副本都在分区却不可用。现在你应该能给出完整解释——要么控制器没能选出 Leader要么 ISR 里没有存活副本要么状态还没从NewPartition或OfflinePartition跃迁到OnlinePartition。这三条分别对应控制面故障、数据面资格不足、状态机未完成。现象最可能的状态优先检查处理方向Leader 为 -1OfflinePartitionISR 是否有存活副本恢复副本或评估 unclean 选举分区一直不可读写NewPartition 卡住Controller 是否稳定修复 Controller 选举ISR 反复收缩扩张副本同步不稳磁盘/网络/负载调 lag 阈值、排查慢副本客户端短时报 NotLeader元数据缓存旧刷新周期与重试属正常切换观察恢复时间选择上可以这样判断当团队已有成熟 ZK 运维且集群规模中等时ZK 模式仍然可行当分区数大、希望减少组件、追求更快的故障转移时选 KRaft。当业务不能容忍丢数据时关闭 unclean 选举、保持副本因子 3 并跨机架分布当业务可用性优先且能接受数据回退时才考虑开启 unclean 选举。最后记住那条主线触发事件 → Controller 决策 → 写元数据 → 通知 Broker → 推进分区状态机 → 选举 Leader → 刷新客户端缓存。任何“分区不可用”的问题都可以沿着这条线一段一段定位。16. 参考资料Apache Kafka 官方文档Replication、Controller、KRaft 相关章节kafka.apache.org/documentation。Apache Kafka 官方文档Topic 配置与 broker 配置参考包含replica.lag.time.max.ms、unclean.leader.election.enable、metadata.max.age.ms。Apache Kafka KIP-500: Replace ZooKeeper with a Self-Managed Metadata Quorum。Apache Kafka KIP-595: A Raft Protocol for the Metadata Quorum。Apache Kafka 源码中PartitionStateMachine、ReplicaStateMachine、KafkaController、PartitionLeaderElectionStrategy相关类。Neha Narkhede 等《Kafka: The Definitive Guide》中关于复制与 Controller 的章节。
返回列表