【发布时间】:2015-12-23 04:02:28
【问题描述】:
下面的设计可以在Storm中完成吗?
让我们以下面的 wordcount 为例 https://github.com/nathanmarz/storm-starter/blob/master/src/jvm/storm/starter/WordCountTopology.java 我正在将单词生成器 spout 更改为文件阅读器 spout
这个字数统计拓扑的设计是 1.Spout读取文件并逐行创建句子 2.螺栓将句子拆分为单词 3. 用螺栓添加唯一词并给出一个词及其对应的计数
因此,拓扑在某种程度上描述了文件需要采用的流程来计算它所具有的唯一词。
如果我有两个文件,文件 1 和文件 2,一个应该能够调用相同的拓扑并创建此拓扑的两个实例以运行相同的字数。
为了跟踪字数统计是否确实已完成,一旦文件处理完毕,字数统计拓扑的实例应处于已完成状态。
在 Storm 的当前设计中,我发现 Topology 是实际的实例,所以它就像一个任务。
一个人需要使用不同的拓扑名称进行两次不同的调用,例如
对于文件 1
StormSubmitter.submitTopology("WordCountTopology1", conf,builder.createTopology());
对于文件 2
StormSubmitter.submitTopology("WordCountTopology2", conf,builder.createTopology());
更不用说使用storm客户端同样上传jar
storm jar stormwordcount-1.0.0-jar-with-dependencies.jar com.company.WordCount1Main.App "server" "filepath1"
storm jar stormwordcount-1.0.0-jar-with-dependencies.jar com.company.WordCount2Main.App "server" "filepath2"
另一个问题是处理文件后拓扑未完成。在我们对拓扑发出 kill 之前,它们一直都处于活动状态
storm kill "WordCountTopology"
我了解,在消息来自消息队列(如 Kafka)的流媒体世界中,消息没有尽头,但在实体/消息已修复的文件世界中,这与它有什么关系。
是否有执行以下操作的 API?
//创建拓扑,这是一次使用storm上传各自的jar
StormSubmitter.submitTopology("WordCountTopology", conf,builder.createTopology());
上传后,应用程序代码只需使用 agruments 实例化拓扑
//创建拓扑实例并提供状态跟踪器
JobTracker tracker = StormSubmitter.runTopology("WordCountTopology", conf, args);
//可以在Storm中查询当前作业是否完成
JobStatus status = StormSubmitter.getTopologyStatus(conf, tracker);
【问题讨论】:
标签: apache-storm