ARTICLE DETAIL

资讯详情

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

图解 Fluss(四):分布式协调 —— 选举、副本状态机与

图解 Fluss(四):分布式协调 —— 选举、副本状态机与 图解 Fluss四分布式协调 —— 选举、副本状态机与 Rebalance阅读本文你将了解三个 Coordinator 如何选出唯一 Leader 且不脑裂、一个 Tablet 副本的状态迁移路径、节点宕机时数据为什么不会丢、以及扩容时 Fluss 如何在不中断服务的前提下搬移数据。配套图表seq-04-startup-leader-election、state-01-replica、state-02-coordinator、seq-05-rebalance难度⭐⭐⭐⭐ | 适合人群负责集群稳定性与容量规划的 SRE / 架构师一、场景大促前的三个运维动作大促前一周你要对 Fluss 集群做三件事动作一冷启动整个集群机房搬迁所有节点关机再开机。开机顺序是什么三个 Coordinator 同时启动谁当 Leader会不会出现两个 Leader动作二应对节点宕机大促期间一台 TabletServer 宕机了。它上面有 2000 个 Tablet其中 300 个是 Leader。这些 Tablet 会发生什么数据会丢吗业务会中断多久动作三紧急扩容流量比预期高 3 倍需要临时加 10 台 TabletServer。加了之后数据会自动均衡过去吗会不会因为大量数据迁移把网络打满反而影响线上业务这三个动作对应的正是本篇的四张图。二、图 1启动与 Leader 选举时序图图里有两个 Coordinator 实例 A 和 B完整演示了选举过程。2.1 阶段 1Standby 模式启动A - A : initCoordinatorStandby() 启动 RPC/元数据/调度器 A - ZK : registerCoordinatorServer() (临时节点) ZK -- A : 注册成功第 1 篇讲过Standby 阶段只起基础设施。关键点是registerCoordinatorServer()创建的是临时节点ephemeral/fluss/coordinator/servers/ ├── server_0000000001 ← A 注册ephemeral ├── server_0000000002 ← B 注册ephemeral └── server_0000000003 ← C 注册ephemeral为什么必须是临时节点因为 ZK 的临时节点在会话断开时会自动删除。Coordinator A 宕机 → 会话超时 → 节点自动消失 → 集群立刻感知到A 不在了。如果用持久节点还需要额外的心跳检测机制。2.2 阶段 2Leader 选举这是整张图最精妙的部分。A - ZK : createEphemeralSequential(/coordinator_election/candidate_) ZK -- A : 节点路径 candidate_001 B - ZK : createEphemeralSequential(/coordinator_election/candidate_) ZK -- B : 节点路径 candidate_002 A - ZK : getChildren(选举路径) ZK -- A : [candidate_001, candidate_002] A - A : 序号 001 最小 → 自己是 Leader B - ZK : getChildren(选举路径) ZK -- B : [candidate_001, candidate_002] B - B : 序号 002 非最小 → Watch candidate_001三个关键设计①ephemeral_sequential临时顺序节点ZK 会自动在节点名后追加一个单调递增的序号/coordinator_election/ ├── candidate_0000000001 ← A ├── candidate_0000000002 ← B └── candidate_0000000003 ← C“序号最小者为 Leader”是一个无需协调的确定性规则——每个节点拿到getChildren()的结果后自己就能判断出谁是 Leader不需要额外的投票轮次。同时因为是 ephemeral 节点Leader 宕机后它的候选节点自动消失下一个序号的节点自然成为新的最小者。② Watch 前一个节点避免羊群效应注意图里 B 的行为B - B : 序号 002 非最小 → Watch candidate_001B watch 的是candidate_001它的前一个节点而不是所有节点也不是 Leader 节点。这个细节叫链式 watch它解决的是羊群效应Herd Effect错误的做法所有非 Leader 节点都 watch Leader 节点 → Leader 宕机时ZK 要向 N-1 个节点同时发送通知 → N-1 个节点同时被唤醒同时发起 getChildren() 查询 → 瞬间产生 N-1 倍的 ZK 读请求 → ZK 被打爆 正确的做法Fluss 采用每个节点只 watch 它前面那一个 → Leader (001) 宕机时ZK 只需要通知 002 → 002 成为新 Leader → 如果 002 也宕机通知 003以此类推 → 通知复杂度 O(1)而不是 O(N)这个模式和 Zookeeper 官方的分布式锁实现Curator 的LeaderLatch是一致的。③ ZK Fence防脑裂A - ZK : fenceBecomeCoordinatorLeader() (epoch 1, 防脑裂) ZK -- A : ZkEpoch(newEpoch)只有成为 Leader 之后才会执行fenceBecomeCoordinatorLeader()把 ZK 上记录的 epoch 加一。图下方的 note 总结了这三点ZooKeeper 选举关键设计: 1. ephemeral_sequential 节点自动处理崩溃退出 2. Watch 前一个节点 (非所有节点) → 避免羊群效应 3. fenceBecomeCoordinatorLeader 用 epoch 防脑裂2.3 阶段 3Leader 初始化A - A : initCoordinatorLeader() 创建 EventProcessor/AutoPartition/ChannelManager A - ZK : registerCoordinatorLeader() (临时节点) A - ZK : createDefaultDatabase(fluss) A - TS : 开始下发 Tablet 分配指令注意registerCoordinatorLeader()创建的是另一个临时节点/fluss/coordinator/leader → {serverId: server_1, epoch: 42, address: ...}这个节点的作用是让TabletServer 和客户端能找到 Leader。TabletServer 启动时 watch 这个节点Leader 变更时立刻感知并重新注册。2.4 回答动作一机房搬迁后冷启动1. 先启动 ZooKeeper 集群3 或 5 节点 ← 必须最先 2. 再启动 3 个 CoordinatorServer ← 自动选主谁先注册谁序号小 3. 最后启动 N 个 TabletServer ← 向 Coordinator 注册等待分配会不会出现两个 Leader不会因为顺序节点的序号由 ZK 保证全局唯一且单调序号最小者为 Leader是确定性规则即使发生网络分区旧 Leader 的 epoch 已过期它下发的指令会被 TabletServer 拒绝三、图 2Coordinator 状态机时序图看的是一次成功的选举过程状态图看的是所有可能的生命周期路径。3.1 四个状态[*] -- Standby : 进程启动 Standby -- Electing : startElectLeaderAsync() Electing -- Leader : 选举成功 (序号最小 ZK Fence) Electing -- Standby : 选举失败 (等待前节点删除) Leader -- Standby : 失去 Leader (会话超时/主动释放) Standby -- Down : closeAsync() Leader -- Down : closeAsync() Down -- [*] : 资源清理完成注意Electing是一个显式的中间状态。这一点很重要——很多人的直觉是启动 → 竞选 → 成功就 Leader但实际上竞选可能失败失败后要回到 Standby 继续等待。3.2 每个状态持有什么资源图里两个 note 精确列出了差异Standby 保留所有节点都有- ZooKeeper 连接 - RpcServer (健康检查) - MetadataManager (只读) - DynamicConfigManager (监听配置)Leader 特有只有一个节点有- CoordinatorEventProcessor (单线程事件循环) - CoordinatorChannelManager (连接 TabletServer) - AutoPartitionManager (自动分区) - RpcClient (主动连接)这个资源划分是第 1 篇两阶段启动的直接体现initCoordinatorStandby()建前者initCoordinatorLeader()建后者。3.3 状态迁移时的资源操作迁移操作Standby → Electing创建顺序节点、getChildren、判断序号Electing → Leaderfence 递增 epoch →initCoordinatorLeader()建 Leader 资源 → 注册 leader 节点Electing → Standbywatch 前一个节点进入等待Leader → StandbycleanupCoordinatorLeader()按逆序清理 Leader 资源 → 重新参与选举图里的 note 特别强调失去 Leader 时按逆序清理这些资源为什么要逆序因为资源之间有依赖创建顺序RpcClient → ChannelManager → AutoPartitionManager → EventProcessor 销毁顺序EventProcessor → AutoPartitionManager → ChannelManager → RpcClient如果先关掉RpcClient正在处理事件的EventProcessor就会因为连接断开而抛异常。逆序销毁保证上层消费者先停下层依赖后关。3.4 代码骨架privatevoidelectCoordinatorLeaderAsync(){coordinatorLeaderElection.startElectLeaderAsync(// 成为 Leader 的回调()-{try{ZkEpochepochzkClient.fenceBecomeCoordinatorLeader();if(epochnull){// fence 失败说明有更新的 epoch放弃return;}initCoordinatorLeader();zkClient.registerCoordinatorLeader();zkClient.createDefaultDatabase(DEFAULT_DATABASE);}catch(Exceptione){// 初始化失败主动释放 Leader 身份重新选举cleanupCoordinatorLeader();}},// 失去 Leader 的回调()-cleanupCoordinatorLeader());}四、图 3Replica 副本状态机Coordinator 管元数据真正的数据可靠性由 Replica 状态机保证。4.1 五个状态[*] -- NewReplica : Tablet 创建 NewReplica : 刚分配尚未加载LogTablet/KvTablet 未打开 OnlineReplica : 已加载数据可服务等待 Leader/Follower 角色 LeaderReplica : 处理读写请求管理 ISR 列表推进 High Watermark FollowerReplica : 从 Leader 拉取日志上报同步进度等待晋升为 Leader OfflineReplica : 暂时不可用可能因故障/网络分区4.2 迁移路径NewReplica -- OnlineReplica : 加载完成 (becomeOnline) OnlineReplica -- LeaderReplica : becomeLeader (Leader 选举成功) OnlineReplica -- FollowerReplica : becomeFollower (跟随 Leader) LeaderReplica -- FollowerReplica : 失去 Leader (被新的 Leader 取代) FollowerReplica -- LeaderReplica : 晋升 (原 Leader 故障) LeaderReplica -- OfflineReplica : 故障 / 网络分区 FollowerReplica -- OfflineReplica : 故障 / 同步超时 OfflineReplica -- OnlineReplica : 恢复 (重新加载) OnlineReplica -- [*] : Tablet 删除注意NewReplica和OnlineReplica的分离副本分配Coordinator 决定和副本加载TabletServer 执行是异步的两步。Tablet 被分配到某台机器不代表它立刻能服务——要等LogTablet打开、索引 mmap、RocksDB 打开完成。4.3 Leader 与 Follower 的职责图里的两个 noteLeader 职责: - 接收客户端写入 (append) - 管理 ISR (加入/移除 Follower) - 推进 HW min(所有 ISR 的 LEO) Follower 职责: - ReplicaFetcher 拉取日志 - 验证 Leader Epoch - 落后太多时通过快照追平验证 Leader Epoch这一条值得展开。这是防止日志分叉log divergence的关键机制。考虑这个场景1. Leader A (epoch5) 写入 offset 100-110还没同步给 Follower 2. A 宕机 3. Follower B 晋升为新 Leader (epoch6)它的日志只到 offset 100 4. B 从 offset 100 开始写入新数据offset 100-105内容不同 5. A 重启它的日志有 offset 100-110 6. A 作为 Follower 向 B 拉取数据 —— 如果不检查A 会保留自己的 100-110 导致同一个 offset 在两个副本上内容不同 → 日志分叉解决办法是Leader Epoch 校验A 作为 Follower 向 B 请求我要从 offset 100 开始拉我的 epoch 是 5 B 查询 epoch 5 对应的最后 offset 是 100返回 OK A 发现自己 offset 100 之后的数据101-110不在有效范围内 → 截断 A 从 B 重新拉取 epoch 6 的数据这个机制和 Kafka 的LeaderEpoch完全一致KIP-101 / KIP-279。4.4 ISR 的动态调整Leader 会持续监控每个 Follower 的同步进度// 伪代码Leader 侧的 ISR 管理voidmaybeShrinkIsr(){longleaderLagMs...;// 配置replica.lag.time.max.ms默认 30sfor(Replicafollower:assignedReplicas){if(now-follower.lastCaughtUpTimeMsleaderLagMs){// Follower 落后超过 30 秒 → 移出 ISRisr.remove(follower);}}}voidmaybeExpandIsr(){for(Replicafollower:assignedReplicas){if(follower.logEndOffsethighWatermark!isr.contains(follower)){// Follower 追上 HW → 重新加入 ISRisr.add(follower);}}}配置replica.lag.time.max.ms的权衡设小如 10s设大如 60s慢副本快速被剔除acksall不受拖累副本容忍抖动不会频繁进出 ISR风险网络偶发抖动导致 ISR 频繁收缩风险真的慢副本拖慢acksall的响应生产建议默认 30s网络环境差的机房可放宽到 60s。4.5 回答动作二一台 TabletServer 宕机上面 2000 个 Tablet300 个 LeaderT0s ZK 会话超时默认 18s节点被摘除 ↓ T18s Coordinator 感知标记这些 Tablet 的副本为 OfflineReplica ↓ T18s 对每个失去 Leader 的 Tablet 从 ISR 中选一个 Follower 晋升为 LeaderReplica 优先选 LEO 最大的保证数据最新 ↓ T19s 新 Leader 开始接收读写 ISR 收缩少了宕机那台的副本 ↓ T?? 如果宕机节点恢复 → NewReplica → OnlineReplica → FollowerReplica → 追平 LEO → 重新加入 ISR数据会丢吗只要写入时用的是acksall且min.insync.replicas 2就不会丢——因为acksall保证数据至少写进了所有 ISR 副本而 ISR 里的副本即使不是 Leader也持有完整数据。业务中断多久主要取决于 ZK 会话超时时间# 缩短故障发现时间但会增加误判风险 zookeeper.session-timeout: 10s replica.lag.time.max.ms: 30s调优提醒把zookeeper.session-timeout设得太小如 3s一次 Full GC 就可能导致误判节点宕机触发不必要的 Leader 切换。生产环境建议不低于 10s。五、图 4Rebalance 渐进式迁移5.1 触发阶段Coord - Rebal : checkAndTriggerRebalance() Rebal - Rebal : 检测负载不均衡 (标准差 阈值) 或新节点加入 Rebal - Rebal : generateRebalancePlan() 计算最优 Tablet 分布 Rebal - Rebal : 生成迁移计划 (最小迁移成本) Rebal -- Coord : RebalancePlan两种触发条件新节点加入集群扩容时自动触发负载不均衡各 TabletServer 的 Tablet 数量标准差超过阈值# 触发阈值配置 tablet-server.rebalance.trigger.threshold: 0.1 # 标准差比例超过 10% 触发 tablet-server.rebalance.check.interval: 5min最小迁移成本是什么Rebalance 规划是一个带约束的优化问题目标minimize(迁移的 Tablet 数 × 每个 Tablet 的数据量) 约束1. 每个 Tablet 的副本不能集中在同一台机器 2. 迁移后各节点的负载标准差 阈值 3. 单个节点同时进行的迁移数 上限5.2 执行阶段渐进式分批图里最关键的结构是这个looploop 每批迁移 (20% Tablet) ... end为什么是渐进式而不是一次性全搬一次性迁移的问题假设要迁移 10000 个 Tablet总计 10TB 数据 - 网络带宽被打满 → 线上业务的读写延迟飙升 - 磁盘 IO 被打满 → 正常读写受影响 - 万一中途出错回滚困难渐进式的好处每批只迁移 20% - 网络和磁盘压力可控 - 每批完成后检查集群稳定性异常就停止 - 支持中途回滚只回滚当前批次5.3 单个 Tablet 迁移的三步图中每个 loop 内部有三个子阶段① 副本同步Src - Dst : 复制 Tablet-X 日志 (作为 Follower) Dst - Src : 追平 Leader (LEO 追上) Dst - ZK : 请求加入 ISR ZK -- Dst : ISR 更新成功注意迁移的第一步是在目标机器上增加一个 Follower 副本而不是直接搬数据。这个顺序保证了任何时刻都有完整的副本可用。② Leader 切换Coord - ZK : 更新 Tablet-X Leader Server-3 ZK -- Coord : Leader 元数据更新 Coord - Dst : 晋升为 Leader Dst -- Coord : 晋升完成因为 Dst 已经追平了 LEO 并在 ISR 里这次切换是干净的——不丢数据业务侧只会感受到毫秒级的抖动和上面节点宕机的切换是同一个机制但更快因为是主动切换。③ 清理旧副本Coord - Src : 删除 Tablet-X 旧副本 Src -- Coord : 删除完成注意清理发生在最后这保证如果第 ② 步失败旧副本还在可以回退。5.4 图下方 note 的总结渐进式迁移三步: 1. Follower 同步 → 追上 Leader 2. 加入 ISR → Leader 切换 3. 删除旧副本 分批执行避免集群压力, 支持中途回滚5.5 回答动作三加 10 台 TabletServer 后T0min 新节点注册到 Coordinator T5min checkAndTriggerRebalance 检测到负载不均衡 生成迁移计划从老节点搬一部分 Tablet 到新节点 T5min 开始第 1 批20% - 每个被迁的 Tablet 先在新节点建 Follower - 追平后切 Leader - 删老副本 T?? 第 1 批完成检查集群稳定性 继续第 2 批 ... T?? 全部完成集群负载均衡会打满网络吗可以通过限流控制# 限制单个 TabletServer 用于副本同步的带宽 tablet-server.replica.fetch.max.bytes: 10MB tablet-server.rebalance.max.concurrent.migrations: 5 # 单节点并发迁移数 tablet-server.rebalance.batch.ratio: 0.2 # 每批比例六、动手验证6.1 观察选举过程# 启动 Coordinator观察日志tail-flogs/coordinator-server.log|grep-EStandby|Elect|Leader# 预期看到# [main] initCoordinatorStandby - Starting coordinator server in standby mode# [main] CoordinatorLeaderElection - Created candidate node: candidate_0000000001# [main] CoordinatorServer - Became coordinator leader with epoch 1# [main] initCoordinatorLeader - Starting coordinator leader services6.2 手动触发 Leader 切换# 找到当前 LeaderzkCli.sh get /fluss/coordinator/leader# kill 掉 Leader 进程kill-9leader_pid# 观察谁接任应该在 10-20 秒内完成zkCli.sh get /fluss/coordinator/leader6.3 观察 Rebalance# 查看当前 Tablet 分布curlhttp://coordinator:9124/metrics|greptablet_distribution# 手动触发 Rebalance如果有管理接口curl-XPOST http://coordinator:9124/rebalance# 观察迁移进度curlhttp://coordinator:9124/metrics|grep-Erebalance|migrat6.4 观察 ISR 变化curlhttp://tablet-server:9125/metrics|grepisr关注指标fluss_replica_isr_count # 每个 Tablet 的 ISR 大小 fluss_replica_under_replicated # 副本数不足的 Tablet 数应该为 0 fluss_replica_leader_count # 每台机器上的 Leader 数应该均衡七、生产实践要点7.1 部署建议项目建议原因Coordinator 数量3 或 5奇数ZK 选举需要多数派ZK 集群独立部署3 或 5 节点混部会因 IO 竞争导致误判ZK session timeout10 - 18s太小易误判太大故障恢复慢TabletServer 数量≥ 3保证副本能分散副本数32 副本在滚动重启时会降级7.2 关键配置# coordinator-server.yaml zookeeper.session-timeout: 15s zookeeper.connection-timeout: 10s coordinator.rebalance.check-interval: 5min coordinator.rebalance.trigger-threshold: 0.1 # tablet-server.yaml replica.lag.time.max.ms: 30000 replica.fetch.max.bytes: 10MB replica.fetch.min.bytes: 1KB replica.fetch.wait.max.ms: 500 # 客户端重要 client.request.acks: all tablet-server.min.insync.replicas: 27.3 滚动重启的正确顺序升级集群时按这个顺序能避免不必要的 Leader 切换1. 先滚动重启 TabletServer一次一台 - 等该节点的 Tablet 全部回到 OnlineReplica 且 ISR 恢复 - 再重启下一台 2. 最后重启 Coordinator - 先重启 Standby 节点 - 最后重启 Leader会触发一次 Leader 切换千万不要先重启 Coordinator Leader——这会导致所有 TabletServer 重新注册产生大量元数据变更事件。7.4 容量规划中的副本开销有效容量 裸容量 / 副本数 / 空间放大系数 举例10 台机器每台 2TB SSD 裸容量 20TB 副本数 3 RocksDB 空间放大 ≈ 1.5PK 表 有效容量 ≈ 20 / 3 / 1.5 ≈ 4.4TBPK 表 ≈ 20 / 3 ≈ 6.7TBLog 表规划时还要留出 30% 余量给 compaction 和 Rebalance 的临时空间。八、排障手册现象可能原因排查方向集群一直无 LeaderZK 不可达或选举路径被污染zkCli.sh ls /fluss/coordinator/election检查是否有残留的candidate_节点Leader 频繁切换GC 停顿 / 网络抖动 / ZK 压力查 GC 日志检查 ZK 的avg latency适当调大 session timeout出现两个 Leaderepoch 相同ZK 脑裂检查 ZK 集群健康度是否过半数存活under_replicated指标持续 0有副本卡在 Offline检查对应节点的磁盘和 GC看replica.lag.time.max.ms是否太小Rebalance 一直不触发标准差未超阈值调小trigger-threshold或手动触发Rebalance 卡住不推进目标节点磁盘满 / 网络限流检查目标节点磁盘水位检查rebalance.max.concurrent.migrations迁移期间业务延迟飙升迁移限速不当调低并发迁移数和replica.fetch.max.bytes节点恢复后一直追不上落后太多需要快照同步看日志有没有 “snapshot” 相关检查网络带宽九、小结四张图串起 Fluss 的分布式协调机制Leader 选举图 1 图 2ephemeral_sequential节点 “序号最小者为 Leader” → 无需投票轮次每个节点只 watch 前一个节点 → 避免羊群效应通知复杂度 O(1)fenceBecomeCoordinatorLeader()递增 epoch → 防止旧 Leader 诈尸两阶段启动Standby 建基础设施Leader 才建协调资源销毁时逆序副本状态机图 3五态New → Online → Leader/Follower → OfflineHW min(ISR 的 LEO)是已提交的边界Leader Epoch 校验防止日志分叉ISR 动态调整滞后超时默认 30s则剔除Rebalance图 4渐进式分批每批 20%避免网络和磁盘被打满单个 Tablet 迁移三步先加 Follower 追平 → 再切 Leader → 最后删旧副本清理在最后保证任何一步失败都能回退三句话记住选举顺序节点定序 链式 watch epoch fence三个设计缺一不可。副本acksallmin.insync.replicas2是不丢数据的底线HW 是已提交边界。迁移先加副本再切主最后删旧任何时候都有完整副本这是不中断服务的关键。下一篇是最后一篇我们把视角拉回到数据怎么进来、怎么演进、怎么归档——Flink Connector、表生命周期与冷热分层。
返回列表