Apache Storm 动态扩缩容:Rebalance、在线调整并行度与拓扑滚动升级实战

发布时间:2026/9/23 12:42:54
Apache Storm 动态扩缩容:Rebalance、在线调整并行度与拓扑滚动升级实战 Apache Storm 动态扩缩容Rebalance、在线调整并行度与拓扑滚动升级实战1. Storm 动态扩缩容概述Apache Storm 是一个开源的分布式实时计算系统广泛应用于流数据处理场景。在生产环境中随着业务负载的变化需要对 Storm 集群进行动态扩缩容以优化资源利用率并保持系统性能。Storm 提供了几种机制来实现动态扩缩容主要包括 Rebalance 命令、在线调整并行度以及拓扑滚动升级。Storm 集群架构与扩缩容组件展示 Storm 集群主要组件及其在动态扩缩容中的作用Nimbus 主节点Supervisor 工作节点Worker 进程Executor 线程ZooKeeper 集群Rebalance 命令并行度调整 API滚动升级工具上图展示了 Storm 集群的典型架构以及实现动态扩缩容的核心组件。Nimbus 主节点负责资源分配和任务调度Supervisor 工作节点托管实际执行计算的 Worker 进程而每个 Worker 包含多个 Executor 线程来处理具体的计算任务。ZooKeeper 作为协调服务存储集群状态和拓扑信息而底部的三种扩缩容机制则提供了不同的动态调整能力。Storm 动态扩缩容的主要优势在于无需停止拓扑即可进行资源调整支持平滑过渡避免服务中断可根据实际负载实时优化资源配置提供细粒度的并行度控制能力2. Rebalance 机制详解Rebalance 是 Storm 提供的一种核心动态扩缩容机制允许在不停止拓扑的情况下重新分配计算资源。通过 Rebalance我们可以动态调整拓扑中各组件的并行度以应对负载变化。Rebalance 工作原理Rebalance 命令通过以下步骤实现资源重新分配向 Nimbus 提交 Rebalance 请求Nimbus 通知所有相关 Supervisor 准备迁移Supervisor 重新分配 Executor 进程重新分配数据流和分组策略Rebalance 机制流程决策树展示 Rebalance 执行过程中的关键决策点和操作流程提交 Rebalance 请求首次执行二次执行重置 Executor 数量?是否有未处理消息?是否是否重置消费位点保持消费位点等待处理完成强制重新分配Nimbus 重新分配资源 → Supervisor 重新启动 Executor → 拓扑恢复运行上图展示了 Rebalance 机制的执行决策树。当提交 Rebalance 请求后系统会判断是首次执行还是二次执行。如果是首次执行系统会询问是否需要重置 Executor 数量如果是二次执行则会检查是否有未处理消息。根据不同的决策路径系统会采取相应操作最终实现资源的重新分配。Rebalance 命令使用示例# 基本语法 storm rebalance topology-name [option1] [option2] ... # 示例将并行度调整为 10等待 30 秒 storm rebalance my-topology -n 10 -w 30 # 示例只重新分配 bolt 的并行度 storm rebalance my-topology -e bolt:5 # 示例强制重新分配丢弃未处理消息 storm rebalance my-topology --forceRebalance 关键参数解释参数描述默认值推荐场景-n,--num-executors设置 Executor 总数保持当前值统一增加/减少资源-e,--executors按组件设置 Executor 数量保持当前值精细化调整特定组件-w,--wait等待时间(秒)0需要优雅停止的情况--force强制重新分配false需要立即生效的场景3. 在线调整并行度实践除了 Rebalance 机制外Storm 还提供了更细粒度的并行度调整能力允许单独调整拓扑中各个组件Spout/Bolt的并行度。并行度调整的工作原理并行度调整通过修改拓扑配置中的parallelism.hint参数实现。当调整并行度后Nimbus 会计算需要的 Executor 数量分配新的资源启动新 Executor 并迁移状态并行度调整流程展示在线调整并行度的详细步骤和状态转换原拓扑配置: spout2, bolt13, bolt22提交并行度调整请求Nimbus 解析新配置计算资源需求新拓扑配置: spout4, bolt15, bolt23 (总 Executor 从 7 增至 12)分配新资源启动新 Executor迁移数据流拓扑继续运行新并行度生效上图展示了并行度调整的完整流程。首先系统会解析新的并行度配置然后计算所需的资源需求分配新资源并启动新的 Executor 进程最后迁移数据流并使新并行度生效。并行度调整实现方法方法一通过 Storm UI 调整登录 Storm Web UI (http://nimbus-host:8080)选择目标拓扑点击 Rebalance 按钮在弹出的对话框中调整并行度参数提交调整请求方法二通过代码动态调整// 创建拓扑配置 TopologyBuilder builder new TopologyBuilder(); // 设置初始并行度 builder.setSpout(my-spout, new MySpout(), 2); builder.setBolt(my-bolt1, new MyBolt1(), 3) .shuffleGrouping(my-spout); builder.setBolt(my-bolt2, new MyBolt2(), 2) .shuffleGrouping(my-bolt1); // 创建配置 Config conf new Config(); conf.setNumWorkers(4); conf.setMaxTaskParallelism(12); // 提交拓扑 StormSubmitter.submitTopology(my-topology, conf, builder.createTopology()); // 后续动态调整并行度的代码 // 获取拓扑配置 TopologyConfig topologyConfig new TopologyConfig(); topologyConfig.setComponentParallelism(my-spout, 4); // 调整 spout 并行度为 4 topologyConfig.setComponentParallelism(my-bolt1, 5); // 调整 bolt1 并行度为 5 topologyConfig.setComponentParallelism(my-bolt2, 3); // 调整 bolt2 并行度为 3 // 提交调整请求 StormSubmitter.rebalance(my-topology, topologyConfig, 30); // 等待 30 秒完成调整方法三通过 Storm CLI 命令调整# 设置拓扑配置文件 cat topology-config.yaml EOF config: spout: parallelism: 4 bolt1: parallelism: 5 bolt2: parallelism: 3 workers: 6 EOF # 应用配置调整 storm rebalance my-topology -f topology-config.yaml并行度调整最佳实践逐步调整避免一次性大幅度调整并行度监控资源调整前后密切关注 CPU、内存使用情况数据倾斜特别关注数据倾斜问题避免部分负载过高状态管理有状态的组件调整并行度需要谨慎处理状态迁移优雅过渡设置适当的等待时间确保数据处理完整4. 拓扑滚动升级策略滚动升级是在不中断服务的情况下更新拓扑代码或配置的重要策略。Storm 提供了多种机制来实现拓扑的滚动升级。滚动升级的工作流程准备新版本的 JAR 包部署新代码到集群逐步升级 Supervisor 节点验证升级结果拓扑滚动升级时间线展示拓扑滚动升级的步骤、时间点和关键操作时间05min10min15min20min准备新版本JAR 包编译部署到共享存储更新代码路径修改 storm.yaml重启 Supervisor触发升级kill -9 旧进程启动新进程验证升级节点检查日志确认指标正常升级下一个节点重复上一阶段逐节点进行全部升级完成监控整体性能回滚准备回滚准备保留旧版本备份记录配置变更问题回滚恢复旧版本恢复原有配置稳定性验证持续监控压力测试上图展示了拓扑滚动升级的时间线从准备新版本到最终完成升级的全过程。整个滚动升级过程可以分为准备阶段、执行阶段和验证阶段每个阶段都有明确的操作步骤和时间点。滚动升级实现方法方法一使用 Storm 原生滚动升级# 1. 准备新版本 JAR cp my-topology-new.jar /shared/storage/ # 2. 更新 storm.yaml 中的 JAR 路径 vim /opt/storm/conf/storm.yaml # 添加或修改以下内容: # topology.jar: /shared/storage/my-topology-new.jar # 3. 重启 Supervisor storm supervisor # 4. 逐步升级拓扑 storm rebalance my-topology --wait 60 # 等待 60 秒让每个节点完成升级 # 5. 验证升级结果 storm list my-topology方法二使用自定义滚动升级脚本#!/bin/bash TOPOLOGY_NAMEmy-topology NEW_JAR_PATH/shared/storage/my-topology-new.jar NODES(node1 node2 node3 node4) WAIT_TIME60 # 1. 部署新 JAR echo Deploying new JAR to shared storage... cp $NEW_JAR_PATH /shared/storage/ # 2. 备份当前配置 echo Backing up current configuration... cp /opt/storm/conf/storm.yaml /opt/storm/conf/storm.yaml.backup # 3. 修改配置 echo Updating storm.yaml... sed -i s|topology.jar:.*|topology.jar: \$NEW_JAR_PATH\| /opt/storm/conf/storm.yaml # 4. 逐节点重启 Supervisor for node in ${NODES[]}; do echo Upgrading node: $node ssh $node sudo systemctl restart supervisor sleep $WAIT_TIME echo Verifying node: $node ssh $node jps | grep -q StormSubmitter done # 5. 验证升级 echo Verifying upgrade... storm list $TOPOLOGY_NAME方法三通过 Storm UI 进行滚动升级登录 Storm Web UI选择目标拓扑点击 Reconfig 按钮上传新的 JAR 包调整配置参数提交升级请求滚动升级注意事项版本兼容性确保新旧版本间的数据格式兼容状态迁移有状态组件升级需要特别处理状态迁移回滚计划提前准备好回滚方案监控告警升级过程中加强监控及时发现异常业务影响评估升级对业务的影响必要时选择低峰期执行实战示例与注意事项下面是一个完整的动态扩缩容实践示例包括 Rebalance、并行度调整和滚动升级的组合使用。完整示例代码import org.apache.storm.Config; import org.apache.storm.LocalCluster; import org.apache.storm.StormSubmitter; import org.apache.storm.generated.StormTopology; import org.apache.storm.topology.TopologyBuilder; import org.apache.storm.tuple.Tuple; import org.apache.storm.tuple.Values; import org.apache.storm.utils.Utils; public class DynamicScalingTopology { public static class SampleSpout extends BaseRichSpout { private SpoutOutputCollector collector; private int count 0; Override public void open(Map conf, TopologyContext context, SpoutOutputCollector collector) { this.collector collector; } Override public void nextTuple() { Utils.sleep(100); collector.emit(new Values(message- count)); } Override public void declareOutputFields(OutputFieldsDeclarer declarer) { declarer.declare(new Fields(message)); } } public static class SampleBolt extends BaseRichBolt { private OutputCollector collector; Override public void prepare(Map stormConf, TopologyContext context, OutputCollector collector) { this.collector collector; } Override public void execute(Tuple input) { String message input.getString(0); System.out.println(Processing: message); // 模拟处理延迟 Utils.sleep(50); collector.ack(input); } Override public void declareOutputFields(OutputFieldsDeclarer declarer) { // 无输出 } } public static StormTopology buildTopology(int spoutParallelism, int boltParallelism) { TopologyBuilder builder new TopologyBuilder(); // 设置初始并行度 builder.setSpout(sample-spout, new SampleSpout(), spoutParallelism); builder.setBolt(sample-bolt, new SampleBolt(), boltParallelism) .shuffleGrouping(sample-spout); return builder.createTopology(); } public static void main(String[] args) throws Exception { Config conf new Config(); conf.setNumWorkers(2); if (args ! null args.length 0) { // 生产环境运行 int spoutParallelism Integer.parseInt(args[0]); int boltParallelism Integer.parseInt(args[1]); StormTopology topology buildTopology(spoutParallelism, boltParallelism); StormSubmitter.submitTopology(dynamic-scaling-topology, conf, topology); } else { // 本地测试 int spoutParallelism 2; int boltParallelism 3; StormTopology topology buildTopology(spoutParallelism, boltParallelism); LocalCluster cluster new LocalCluster(); cluster.submitTopology(dynamic-scaling-topology, conf, topology); // 模拟扩缩容 Thread.sleep(10000); System.out.println(Increasing parallelism...); cluster.rebalance(dynamic-scaling-topology, 30000, new TopologyConfig().setComponentParallelism(sample-spout, 4) .setComponentParallelism(sample-bolt, 5)); Thread.sleep(30000); System.out.println(Decreasing parallelism...); cluster.rebalance(dynamic-scaling-topology, 30000, new TopologyConfig().setComponentParallelism(sample-spout, 2) .setComponentParallelism(sample-bolt, 3)); Thread.sleep(30000); cluster.shutdown(); } } }最小运行示例编译打包:mvn clean package -DskipTests启动拓扑:storm jar target/dynamic-scaling.jar DynamicScalingTopology 2 3动态扩容:storm rebalance dynamic-scaling-topology -e sample-spout:4,sample-bolt:5 -w 30动态缩容:storm rebalance dynamic-scaling-topology -e sample-spout:2,sample-bolt:3 -w 30注意事项数据一致性扩缩容过程中可能出现数据重复或丢失需要设计适当的去重或容错机制资源争用调整并行度时注意系统资源争用问题避免过度占用导致集群不稳定监控告警建立完善的监控体系及时发现扩缩容过程中的异常情况渐进式调整避免一次性大幅调整并行度采用渐进式调整策略业务影响评估评估扩缩容对业务的影响选择合适的调整时机通过合理运用 Rebalance、在线调整并行度和拓扑滚动升级等技术可以有效地实现 Storm 集群的动态扩缩容提高资源利用率和系统稳定性满足不同业务场景下的需求变化。