【发布时间】:2013-03-24 03:12:03
【问题描述】:
我正在尝试将 Storm (see here) 集成到我的项目中。我理解了拓扑、spout 和 bolts 的概念。但现在,我试图弄清楚一些事情的实际实现。
A) 我有一个使用 Java 和 Clojure 的多语言环境。我的 Java 代码是一个回调类,其中包含触发流数据的方法。推送到这些方法的事件数据是我想用作 spout 的。
所以第一个问题是如何将进入这些方法的数据连接到一个 spout ?我正在尝试 i) 传递一个 backtype.storm.topology.IRichSpout ,然后 ii) 传递一个 backtype.storm。 spout.SpoutOutputCollector (see here) 到该 spout 的 open 函数 (see here)。但我看不到实际传递任何类型的地图或列表的方法。
B) 我项目的其余部分都是 Clojure。通过这些方法将有大量数据。每个事件的 ID 介于 1 和 100 之间。在 Clojure 中,我希望将来自 spout 的数据拆分到不同的执行线程中。我认为,这些将是螺栓。
如何设置 Clojure bolt 以从 spout 获取事件数据,然后根据传入事件的 ID 中断线程?
提前致谢 蒂姆
[编辑 1]
我实际上已经解决了这个问题。我最终1)实现了我自己的 IRichSpout。然后我 2) 将该 spout 的内部元组连接到我的 java 回调类中的传入流数据。我不确定这是否是惯用的。但它编译并运行没有错误。但是,3) 我没有看到通过 printstuff 螺栓传入的流数据(肯定在那里)。
为了确保事件数据得到传播,我在 spout 或 bolt 实现或拓扑定义中是否需要做一些具体的事情?谢谢。
;;将 Java 回调绑定到我创建的 Spout (.setSpout java-callback ibspout) (storm/defbolt printstuff ["word"] [元组收集器] (println (str "printstuff --> tuple["tuple"] > collector["collector"]")) ) (风暴/拓扑 { "1" (storm/spout-spec ibspout) } { "3" (storm/bolt-spec { "1" :shuffle } 印刷品 ) })[编辑 2]
根据 SO 成员 Ankur 的建议,我正在调整我的拓扑结构。创建 Java 回调后,我使用 (.setTuple ibspout (.getTuple java-callback)) 将它的元组传递给下面的 IBSpout。我没有传递整个 Java 回调对象,因为我得到了 NotSerializable 错误。一切都编译并运行没有错误。但同样,我的 printstuff 螺栓没有数据。嗯。
【问题讨论】:
标签: clojure streaming message-queue apache-storm