ARTICLE DETAIL

资讯详情

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

Apache Storm Nimbus 高可用(HA)设计深度解析:主备选举、代码分发与配置实战

Apache Storm Nimbus 高可用(HA)设计深度解析:主备选举、代码分发与配置实战 大数据流处理后端【免费下载链接】stormApache Storm项目地址https://gitcode.com/gh_mirrors/storm6/storm点击查看免费下载本篇技术指南围绕 Apache Storm 的 Nimbus 高可用设计文档docs/nimbus-ha-design.md展开系统讲解 Nimbus 主备模式的 Leader 选举机制、拓扑代码在多个 Nimbus 主机间的复制分发协议、ClusterSummary/NimbusSummaryThrift API 扩展以及topology.min.replication.count等核心配置的完整参数语义。读完本文你将掌握如何把单点 Nimbus 集群改造成具备故障自动接管能力的 HA 集群理解ILeaderElector、代码分发接口与 BlobStore 复制路径的底层实现并能在生产环境合理调优相关配置。一、问题背景为什么 Nimbus 需要高可用在默认部署中Storm 的 master即 Nimbus是运行在单台机器上、受 supervisor 监控的进程。大多数情况下 Nimbus 的故障是瞬时性的可以被自动重启恢复但一旦出现磁盘损坏、网络分区network partition等严重故障Nimbus 就会长时间不可用。此时集群会出现如下症状已有的拓扑仍然正常运行worker 与 supervisor 之间的数据面不受 Nimbus 影响但无法提交新拓扑也无法对已有拓扑执行kill / deactivate / activate等管理操作若某个 supervisor 节点同时故障则调度重分配reassignment无法执行导致性能退化甚至拓扑失败。HA 项目的目标正是解决这一单点问题让 Nimbus 运行在primary-backup主备模式下保证即使某个 Nimbus 服务器宕机也有备份节点接管其职责。设计文档还明确了四条硬性需求docs/nimbus-ha-design.md提高 Nimbus 的总体可用性increase overall availability允许 Nimbus 主机随时离开或加入集群新加入的主机应能自动追赶进度auto catch up并进入潜在 leader 列表Nimbus 故障切换failover后无需重新提交拓扑任何活跃拓扑都不允许丢失。二、Leader 选举从设计接口到 Zookeeper 实现2.1 设计文档中的ILeaderElector接口设计文档给出了一套面向选主的抽象接口核心语义如下原文摘录public interface ILeaderElector { void addToLeaderLockQueue(); // 排队竞争领导锁立即返回调用者需自行检查 isLeader() void removeFromLeaderLockQueue(); // 退出队列若当前是 leader 则同时释放锁 boolean isLeader(); // 当前调用者是否持有领导锁 InetSocketAddress getLeaderAddress(); // 获取当前 leader 地址无人持锁时抛异常 ListInetSocketAddress getAllNimbusAddresses(); // 当前所有 Nimbus 地址含 leader }2.2 仓库中的实际接口演进当前仓库中该接口位于 storm-client/src/jvm/org/apache/storm/nimbus/ILeaderElector.java相比设计文档已演进得更为完善ILeaderElector extends AutoCloseable并新增了以下关键方法void prepare(MapString, Object conf)初始化阶段保证被调用的配置入口void quitElectionFor(int delayMs)退出竞选若是 leader 则让出领导权并在指定延迟后重新排队这是文档中removeFromLeaderLockQueue的强化版用于处理代码未同步齐、暂让领导权的场景NimbusInfo getLeader()返回 leader 地址无人持锁时可返回null注意NimbusInfo取代了文档中的InetSocketAddress见 storm-client/src/jvm/org/apache/storm/nimbus/NimbusInfo.javaboolean awaitLeadership(long timeout, TimeUnit timeUnit)带VisibleForTesting注解仅供单 Nimbus 集群下 LocalCluster 测试等待 Nimbus 获得领导权后提交拓扑ListNimbusInfo getAllNimbuses()返回当前所有 Nimbus 地址列表含 leader。接口还明确要求addToLeaderLockQueue()具备幂等性可被多次调用并且调用者在失去锁或排队位置丢失时如 Zookeeper 会话重置实现方必须更新内部状态使isLeader()始终反映真实状态。2.3 Zookeeper 实现Curator LeaderLatch 与/leader-lock第一个落地实现是基于 Zookeeper 的当前仓库的实现在 storm-server/src/main/java/org/apache/storm/zookeeper/LeaderElectorImp.java核心要点使用 Curator 框架的LeaderLatch竞逐领导权锁节点路径固定为/leader-lock每个 Nimbus 以自身节点 id 作为LeaderLatch的参与标识addToLeaderLockQueue()内部处理三种 latch 状态CLOSED状态会重建 latch 并重新排队LATENT未启动状态会注册 leader 监听器并start()已排队则直接记录日志——这保证了方法的幂等性通过LeaderListenerCallbackFactorystorm-server/src/main/java/org/apache/storm/zookeeper/LeaderListenerCallbackFactory.java为 latch 注册状态变化监听当节点当选/落选 leader 时回调 Nimbus触发后续的领导权交接动作内部使用StormTimer名为leader-elector-timer调度quitElectionFor的延迟重新入队逻辑。2.4 领导权接管前的完整性检查与让位协议设计文档规定了一条关键不变量leader Nimbus 必须在本地持有所有活跃拓扑的代码。具体流程为Nimbus 启动后先检查本地是否拥有全部活跃拓扑/storm/storms/下登记的拓扑的代码就绪后才调用addToLeaderLockQueue()排队竞争领导权被通知当选 leader 时再次检查本地代码完整性——若缺失任何活跃拓扑代码则拒绝接受领导权释放锁待补全代码后重新入队。当前仓库在 Nimbus.java 中可以看到完整的配套逻辑launchServer()中先通过state.addNimbusHost(host, nimbusSummary)将自身登记进 Zookeeper 的 nimbus 列表再调用leaderElector.addToLeaderLockQueue()入队选主L1528-L1531随后用timer.scheduleRecurring(3, 5, ...)每 5 秒检查一次isLeader()一旦检测到本节点成为 leader 且此前不是 leader便对 Zookeeper 中所有活跃拓扑执行TopologyActions.GAIN_LEADERSHIP状态迁移完成领导权交接L1541-L1555。设计注释明确指出之前一次性检查的实现方式会漏掉GAIN_LEADERSHIP转换因此改为周期性轮询这也解决了 HA Nimbus 新当选后的交接遗漏问题诸如doCleanup()清理拓扑与 jar这类只有 leader 才能执行的操作入口处都先做if (!isLeader()) return;的守卫判断L3155-L3158。对于非 leader 节点收到只有 leader 才能执行的请求文档规定应抛出RuntimeException。仓库中的实现是assertIsLeader()L1919-L1924private void assertIsLeader() throws Exception { if (!isLeader()) { NimbusInfo leaderAddress leaderElector.getLeader(); throw new RuntimeException(not a leader, current leader is leaderAddress); } }该守卫在提交拓扑、设置 blob 复制、调度等多个只有 leader 能执行的入口被调用如 L2007、L2111、L3559、L4459从源码结构看这正是文档所述非 leader 收到请求即抛运行时异常约定的具体落地。2.5 一次完整的 Nimbus 故障切换Failover演练设计文档用一组具体数字描述了故障切换过程这是理解整个协议最好的例子假设集群运行4 个拓扑、共3 个 Nimbus 节点code-replication-factor 2即每个拓扑的代码至少复制到 2 台主机初始不变量成立leader 本地持有全部 4 个拓扑的代码nonLeader-1 持有前 2 个拓扑的代码nonLeader-2 持有后 2 个拓扑的代码leader 宕机且是硬盘故障、无法恢复nonLeader-1 收到 Zookeeper 通知成为新 leader但在接受领导权前检查发现本地只有 2 个拓扑的代码于是放弃锁relinquish转而在/storm/code-distributor/topologyId下查询可下载代码/元文件的来源它找到 leader 与 nonLeader-2 两个条目并通过重试机制尝试从两者下载缺失代码nonLeader-2 的代码同步后台线程也发现自己缺失 2 个拓扑的代码执行同样的下载流程最终至少一个 Nimbus 补齐全部代码并接受领导权集群恢复正常管理能力。整个 leader 选举与故障切换的组件交互时序如下原文档附图三、Nimbus 状态存储与代码分发Code Distribution3.1 Nimbus 的两类状态数据Nimbus 持久化两类数据元信息supervisor 信息、assignment 调度信息等存放在 Zookeeper 中拓扑实际代码拓扑配置与 jar存放在 Nimbus 主机的本地磁盘上。要实现主备切换Nimbus 的状态/数据必须复制到所有 Nimbus 主机或存入分布式存储。文档明确权衡了两条路线精确的数据复制涉及状态管理、一致性校验且正确性很难测试而许多 Storm 用户不愿意为高可用额外引入 HDFS 之类的复制型存储依赖。因此设计上优先采用基于本地文件系统的代码分发 后台同步方案并预留了演进空间考虑到 jar 的体积与超大规模 supervisor 集群的扩展性远期计划迁移到 BitTorrent 协议做代码分发——因为文件系统式的分发模型无法支撑非文件系统式的分发如 BitTorrent。3.2ICodeDistributor接口设计为同时支持 BitTorrent 与各类文件系统复制存储设计文档提出了代码分发抽象public interface ICodeDistributor { void prepare(Map conf); // 初始化 File upload(Path dirPath, String topologyId); // 上传本地代码目录返回 Meta 文件 ListFile download(Path destDirPath, String topologyid, File metafile); // 依 Meta 文件下载代码 int getReplicationCount(String topologyId); // 返回代码已被复制到的主机数 void cleanup(String topologyid); // 清理 void close(Map conf); // 关闭 }其中upload返回的Meta 文件必须包含足够的信息供下载方定位代码对 BitTorrent 而言它是 torrent 文件对 HDFS/S3 而言它记录实际目录或待下载文件的路径。需要说明的是该接口仅存在于设计文档中当前仓库源码内已检索不到ICodeDistributor接口与其默认实现LocalFileSystemCodeDistributor的类文件同时 conf/defaults.yaml L64 仍保留storm.codedistributor.class配置项值为org.apache.storm.codedistributor.LocalFileSystemCodeDistributor。从当前 Nimbus.java 的实现看这一设计已演进为基于 BlobStore 的复制路径详见 3.4 与第四节阅读旧版设计文档时需留意这一差异。3.3 复制一致性的权威来源引入复制必然带来一致性问题。设计文档给出的约定非常清晰以 Zookeeper 中登记的活跃拓扑列表作为代码必须存在的权威authority。据此派生出一条核心规则任何 Nimbus 主机只要缺少 Zookeeper 中标记为 active 的任一拓扑代码就必须让出领导锁把成为 leader 的机会让给其他节点与此同时所有 Nimbus 主机上的后台同步线程持续从代码已成功复制的节点拉取缺失代码因此只要每个活跃拓扑至少存在一个 seed 主机最终就至少有一个 Nimbus 能补齐代码并接受领导权。3.4 拓扑代码复制Replication的完整流程设计文档描述了提交一个拓扑时的代码复制时序客户端上传 jar行为与原来完全一致无任何变化客户端提交拓扑leader Nimbus 调用 code distributor 的upload函数生成 Meta 文件并本地存储随后在/storm/code-distributor/topologyId下写入新条目通知所有非 leader Nimbus 下载新代码leader Nimbus等待至少 N 个非 leader 节点完成代码复制N 由用户配置等待时间有用户可配置的超时非 leader Nimbus 收到新代码通知后先从 leader 下载 Meta 文件再调用 code distributor 的download函数以 Meta 文件为输入下载真实代码下载完成后该非 leader 节点在/storm/code-distributor/topologyId下登记自己成为 leader 宕机时的备用代码来源最后 leader Nimbus 继续执行提交拓扑的常规流程创建 assignment 等。拓扑提交场景下各组件通信的时序如下图所示原文档附图3.5 当前仓库的 BlobStore 复制实现在 Nimbus.java 中等价的等待复制完成逻辑由waitForDesiredCodeReplication(topoConf, topoId)L2089-L2126承担。它围绕 BlobStore 中的三个 key 轮询复制计数jarConfigUtils.masterStormJarKey(topoId)本地模式下跳过 jar 检查代码ConfigUtils.masterStormCodeKey(topoId)配置ConfigUtils.masterStormConfKey(topoId)。getBlobReplicationCount(key)L2081-L2087通过blobStore.getBlobReplication(key, NIMBUS_SUBJECT)读取当前复制数。主循环在三个计数均未达到topology.min.replication.count时每秒重查一次并在每个周期调用assertIsLeader()校验自己仍是 leader防止等待期间领导权已转移一旦累计等待超过topology.max.replication.wait.time.sec且该值为正则打印告警日志并放行继续执行拓扑激活。这与设计文档中等待 N 个节点复制完成 用户可配置超时的语义完全对应。另外blobstore 自身的复制因子由storm.blobstore.replication.factor控制默认值为 3见 conf/defaults.yaml L180与 Nimbus HA 的拓扑代码复制是两个层次的机制前者是 BlobStore 内部的数据冗余后者是拓扑 jar/config/code 在多个 Nimbus 主机间的分发冗余。四、Thrift 与 REST API让客户端发现 leader为避免 worker / supervisor / UI 各自直连 Zookeeper 去查询 master Nimbus 地址设计文档提出改造getClusterInfoAPI它原本返回的ClusterSummary只含supervisorSummary与topologySummary列表现在新增NimbusSummary列表struct ClusterSummary { 1: required listSupervisorSummary supervisors; 3: required listTopologySummary topologies; 4: required listNimbusSummary nimbuses; } struct NimbusSummary { 1: required string host; 2: required i32 port; 3: required i32 uptime_secs; 4: required bool isLeader; 5: required string version; }该结构已落地到当前仓库的生成代码 storm-client/src/jvm/org/apache/storm/generated/ClusterSummary.javaClusterSummary持有ListNimbusSummary nimbuses字段L40并提供了add_to_nimbuses、get_nimbuses_iterator等访问方法同时 Nimbus 启动时通过new NimbusSummary(host, port, uptimeSecs, false, STORM_VERSION)将自己登记到 ZookeeperNimbus.java L1528-L1530isLeader初始为false当选后由选主逻辑更新。这套 API 的使用方包括StormSubmitter、Nimbus 客户端、supervisor 与 UI它们借此发现当前 leader 与所有参与的 Nimbus 主机且任何一台 Nimbus 主机都能响应该类请求。为了降低 Zookeeper 压力文档给出的实践是Nimbus 从 Zookeeper 读取一次该信息并缓存仅在 watcher 被触发选举变化正常情况下极少发生时更新缓存。五、HA 相关配置详解以当前仓库默认值为准Nimbus HA 开箱即用但默认配置只假定单个 Nimbus 主机其取舍是牺牲复制换取更低的拓扑提交延迟。生产环境可按需调整以下参数默认值以 conf/defaults.yaml 为准配置项作用文档描述默认值当前仓库默认值storm.codedistributor.class代码分发器实现类全限定名需实现ICodeDistributororg.apache.storm.codedistributor.LocalFileSystemCodeDistributor与文档一致L64topology.min.replication.count拓扑被标记为 active 并生成 assignment 前代码必须复制到的最少 Nimbus 主机数11L103topology.max.replication.wait.time.sec等待复制达到min.replication.count的最长时间超时后即使未达标也继续执行激活-1 表示永远等待60 秒60L104nimbus.code.sync.freq.secsNimbus 后台线程同步本地缺失拓扑代码的频率5 分钟120 秒L89关于storm.codedistributor.class设计文档给出两个候选LocalFileSystemCodeDistributor默认用本地文件系统同时存放 Meta 文件与代码/配置。缺点是对 Zookeeper 有额外负担——即使下载完 code-distributor 的 Meta 文件仍需联系 Zookeeper 确定可下载真实代码/配置的主机、并查询当前复制计数HDFSCodeDistributor文档中类名为org.apache.storm.hdfs.ha.codedistributor.HDFSCodeDistributor依赖 HDFS不额外加重 Zookeeper 负担拓扑提交更快。注意当前仓库源码中检索不到该实现类它只出现在设计文档中是否可用请以你实际使用的 Storm 发行版或独立的 HDFS 模块为准。关于nimbus.code.sync.freq.secs需特别提示设计文档写作时默认值是 5 分钟而当前仓库 conf/defaults.yaml L89 的默认值已是120 秒配置时请以当前版本实际生效值为准。5.1 一个重要的工程经验提交延迟与后台同步频率的关系设计文档特别指出了一条实践中观察到的规律务必重视尽管所有 Nimbus 主机都注册了 Zookeeper watcher理论上新拓扑代码一出现就会立即触发下载但实际回调几乎不会真的触发代码下载。实践中期望的复制desired replication只有在后台同步线程运行后才会达成。因此当topology.min.replication.count 1时拓扑提交耗时大致处于0 到 (2 × nimbus.code.sync.freq.secs)之间。结合 Nimbus.java 中waitForDesiredCodeReplication的轮询实现可以理解leader 会在超时窗口内每秒检查复制计数复制是否真的达成取决于后台同步线程的调度节奏而不是 watcher 的即时性。这意味着若你的拓扑提交对延迟敏感应适当调小nimbus.code.sync.freq.secs当前默认 120 秒意味着提交最坏可能多等约 4 分钟但频率过高会加重 Nimbus 主机的磁盘 I/O 与网络负载需结合实际集群规模权衡若集群只有一个 Nimbustopology.min.replication.count 1复制等待几乎立即可满足提交延迟基本不受影响。六、深入阅读指引与验证路径如果希望进一步验证本文涉及的设计与实现可以按以下路径在仓库中继续深入设计文档原文docs/nimbus-ha-design.md选主接口storm-client/src/jvm/org/apache/storm/nimbus/ILeaderElector.javaZookeeper 选主实现storm-server/src/main/java/org/apache/storm/zookeeper/LeaderElectorImp.java以及监听器工厂 storm-server/src/main/java/org/apache/storm/zookeeper/LeaderListenerCallbackFactory.javaNimbus 端 HA 集成storm-server/src/main/java/org/apache/storm/daemon/nimbus/Nimbus.java 中的launchServer()L1519 起、assertIsLeader()L1919、waitForDesiredCodeReplication()L2089Nimbus 地址信息模型storm-client/src/jvm/org/apache/storm/nimbus/NimbusInfo.java支持host:port与host:port:tlsPort两种解析格式Thrift 生成结构storm-client/src/jvm/org/apache/storm/generated/ClusterSummary.java全部相关配置默认值conf/defaults.yamlstorm.codedistributor.classL64、nimbus.code.sync.freq.secsL89、topology.min.replication.countL103、topology.max.replication.wait.time.secL104、storm.blobstore.replication.factorL180。七、小结Nimbus 高可用通过主备模式 领导权可移交 代码多机复制三条主线把原先的单点故障转化为可自动恢复的常态故障ILeaderElector抽象了选主协议Zookeeper 实现基于 Curator LeaderLatch 与/leader-lock锁路径接受领导权前必须持有全部活跃拓扑代码的完整性检查与让位机制保证了任何时刻集群管理能力可无缝交接代码分发层以 Zookeeper 活跃拓扑列表为一致性权威配合后台同步线程与复制计数等待逻辑确保即使 leader 硬盘级故障其余节点也能补齐代码完成接管ClusterSummary/NimbusSummary的 API 扩展则让客户端、supervisor 与 UI 无需直连 Zookeeper 即可发现当前 leader。部署 HA 集群时只需按第四节给出的参数语义调整topology.min.replication.count、topology.max.replication.wait.time.sec与nimbus.code.sync.freq.secs并理解提交延迟与后台同步频率之间的权衡即可。赞分享大数据流处理后端【免费下载链接】stormApache Storm项目地址https://gitcode.com/gh_mirrors/storm6/storm点击查看免费下载相关推荐NSwag实战指南3步构建全栈API开发流水线NSwag实战指南3步构建全栈API开发流水线 NSwag作为.NET生态中强大的Swagger/OpenAPI工具链能够高效连接后端API与前端应用实现开发工具代码生成API设计Apache Storm高可用架构Nimbus HA与故障转移机制终极指南Apache Storm高可用架构Nimbus HA与故障转移机制终极指南 Apache Storm作为分布式实时计算系统其高可用架构是保障数据处理连续性的流处理后端大数据Apache Storm 容错机制深度解析从 Worker 失效到 Nimbus 高可用与数据不丢失保证Apache Storm 容错机制深度解析从 Worker 失效到 Nimbus 高可用与数据不丢失保证 Apache Storm 作为流式处理系统其核心卖大数据流处理后端上一篇Hermes-2-Pro-Mistral-7B-SFT代码生成能力详解从简单函数到复杂项目下一篇MOSS-Music-8B-Instruct应用场景大全从音乐教育到内容创作的10个实际用例创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表