【问题标题】:Storm > Howto Integrate Java callback into a SpoutStorm > 如何将 Java 回调集成到 Spout 中
【发布时间】: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 螺栓没有数据。嗯。

公共类 IBSpout 实现 IRichSpout { /** *风暴喷口的东西 */ 私人 SpoutOutputCollector _collector; 私有列表 _tuple = new ArrayList(); public void setTuple(List tuple) { _tuple = tuple; } 公共列表 getTuple() { return _tuple; } /** * Storm ISpout 接口函数 */ public void open(Map conf, TopologyContext context, SpoutOutputCollector 收集器) { _collector = 收集器; } 公共无效关闭(){} 公共无效激活(){} 公共无效停用(){} 公共无效 nextTuple() { _collector.emit(_tuple); } 公共无效确认(对象 msgId){} 公共无效失败(对象 msgId){} public void declareOutputFields(OutputFieldsDeclarer 声明者) {} public java.util.Map getComponentConfiguration() { return new HashMap(); } }

【问题讨论】:

    标签: clojure streaming message-queue apache-storm


    【解决方案1】:

    您似乎将 spout 传递给您的回调类,这似乎有点奇怪。当一个拓扑被执行时,storm 会定期调用 spouts nextTuple 方法,因此您需要做的是将 java 回调传递给您的自定义 spout 实现,以便当storm 调用您的 spout 时,spout 调用 java 回调以获取下一个一组要输入拓扑的元组。

    要理解的关键概念是 Spouts 在风暴请求时数据,您不要将数据推送到 spouts。您的回调不能调用 spout 将数据推送给它,而是当您的 spout 的 nextTuple 方法被调用时,您的 spout 应该从一些 java 方法或任何内存缓冲区中提取数据。

    【讨论】:

    • 哦,太好了。感谢您的洞察力。但是我仍然没有看到数据通过喷口传到我的螺栓上。我在上面给出了更好的描述。也许我应该对我的 Spout 做一些特定的事情?是否有一种特殊的方式必须将数据结构传递给 Spout?谢谢。
    • @Nutritioustim 你得到答案了吗?
    • 嘿嘿。见上文。我无法得到暴风雨来做我想做的事。 Lamina 是一个更轻量级的工具,可以解决我的问题。 HTH。
    【解决方案2】:

    B部分答案:

    在我看来,直截了当的答案就像您正在寻找一个字段分组,这样您就可以控制在执行期间按 ID 将哪些作品分组在一起。

    也就是说,我不确定这是否真的是一个完整的答案,因为我不知道您为什么要这样做。如果您只想要平衡的工作负载,则随机分组是更好的选择。

    【讨论】:

    • 嘿,谢谢你看这个。我实际上确实指定了一个 :shuffle 来平衡工作量。我现在遇到的问题是我没有看到我的事件数据传播到我的螺栓(参见上面的编辑)。任何见解,都值得赞赏。
    • @Nutritioustim 你真的找出问题所在了吗?
    • @Vor,不。Storm 对于我正在尝试做的事情来说似乎有点太不可行了。现在 Lamina 满足了我的需求。 HTH。
    • 谢谢回复,我去看看
    猜你喜欢
    • 2018-09-19
    • 2015-12-28
    • 2013-06-24
    • 2017-04-09
    • 1970-01-01
    • 1970-01-01
    • 2018-11-03
    • 1970-01-01
    • 2014-02-11
    相关资源
    最近更新 更多