如何让Storm的submitTopology执行完成后再终止拓扑并关闭集群?
我正在阅读《Storm应用实战》一书,书中有如下代码片段:
LocalCluster lc = new LocalCluster(); lc.submitTopology("GitHub-commit-count-topology", config, topology); Utils.sleep(TEN_MINUTES); lc.killTopology("GitHub-commit-count-topology"); lc.shutdown();
这段代码会提交拓扑执行,等待固定10分钟后终止拓扑并关闭集群,但这种方式并不合理。我希望实现提交拓扑后等待其执行完成,再执行终止和关闭操作,就像Akka Streams中通过等待Future[Done]完成的方式一样,请问该如何实现?
嗨,这个问题我之前也踩过坑——固定时长等待真的太鸡肋了,要么白等半天浪费资源,要么提前停了导致任务没做完。Storm本身确实没有像Akka Streams那样直接给个Future[Done]的现成API,但咱们可以用几种方式实现“等拓扑干完活再收摊”的逻辑:
方法1:基于拓扑状态的轮询检查
如果你的拓扑是有限流处理(比如处理完一批固定的数据就结束),Storm的拓扑会在所有Spout都完成发送、所有Tuple都被处理ack之后进入COMPLETED状态。我们可以通过轮询检查拓扑状态来实现等待:
LocalCluster lc = new LocalCluster(); lc.submitTopology("GitHub-commit-count-topology", config, topology); // 轮询检查拓扑状态,每隔10秒查一次 while (true) { TopologyStatus status = lc.getTopologyStatus("GitHub-commit-count-topology"); if (status.getStatus() == TopologyStatus.Status.COMPLETED) { break; } Utils.sleep(10000); } // 拓扑完成后再执行关闭操作 lc.killTopology("GitHub-commit-count-topology"); lc.shutdown();
⚠️ 注意:这个方法只适用于有限流拓扑,如果是无限流(比如持续接收实时数据),拓扑永远不会进入COMPLETED状态,这种情况你需要自定义结束条件。
方法2:自定义完成信号(通用场景)
如果你的拓扑有明确的结束触发条件(比如处理完N条数据、某个时间点到达),可以在拓扑内部设置一个线程安全的完成标记,然后在外部等待这个标记被触发:
- 先在Spout/Bolt里维护结束标记,比如在发送完所有数据的Spout中:
public class FiniteSpout extends BaseRichSpout { private final AtomicBoolean isCompleted; private int totalTuples; private int emittedCount = 0; // 通过构造函数传入外部的标记 public FiniteSpout(AtomicBoolean isCompleted) { this.isCompleted = isCompleted; } @Override public void open(Map conf, TopologyContext context, SpoutOutputCollector collector) { totalTuples = (int) conf.get("total.tuples"); } @Override public void nextTuple() { if (emittedCount < totalTuples) { // 发送业务数据 collector.emit(new Values("commit-data-" + emittedCount)); emittedCount++; } else { // 所有数据发送完成,标记状态 isCompleted.set(true); // 休眠避免空轮询占用资源 Utils.sleep(1000); } } // 其他方法省略... }
- 然后在外部创建标记、传递给拓扑,再等待标记触发:
AtomicBoolean topologyCompleted = new AtomicBoolean(false); // 构建拓扑时传递标记 TopologyBuilder builder = new TopologyBuilder(); builder.setSpout("github-commit-spout", new FiniteSpout(topologyCompleted), 1); // 添加Bolt等其他组件... LocalCluster lc = new LocalCluster(); lc.submitTopology("GitHub-commit-count-topology", config, builder.createTopology()); // 等待标记变为true while (!topologyCompleted.get()) { Utils.sleep(5000); } // 执行关闭操作 lc.killTopology("GitHub-commit-count-topology"); lc.shutdown();
这种方式不管是有限流还是需要自定义结束条件的无限流场景都适用,灵活性很高。
方法3:用TopologyListener实现事件驱动等待
Storm的LocalCluster支持添加TopologyListener,可以监听拓扑的状态变化事件,当拓扑进入COMPLETED状态时自动触发关闭逻辑,不需要轮询,更优雅:
LocalCluster lc = new LocalCluster(); // 用CountDownLatch来实现等待逻辑 CountDownLatch completionLatch = new CountDownLatch(1); // 添加拓扑状态监听器 lc.addTopologyListener(new TopologyListener() { @Override public void topologyStarted(String topologyId) {} @Override public void topologyShutdown(String topologyId) {} @Override public void topologyRemoved(String topologyId) {} @Override public void topologyStatusChanged(String topologyId, TopologyStatus.Status newStatus) { // 当目标拓扑进入COMPLETED状态时,触发latch if (newStatus == TopologyStatus.Status.COMPLETED && topologyId.equals("GitHub-commit-count-topology")) { completionLatch.countDown(); } } }); lc.submitTopology("GitHub-commit-count-topology", config, topology); // 等待latch被触发,中断时恢复线程中断状态 try { completionLatch.await(); } catch (InterruptedException e) { Thread.currentThread().interrupt(); } // 关闭拓扑和集群 lc.killTopology("GitHub-commit-count-topology"); lc.shutdown();
这个方法适合有限流拓扑,代码更简洁,不需要自己写轮询逻辑。
内容的提问来源于stack exchange,提问作者Knows Not Much

