【问题标题】:Multiple Micro batch Storm Topology多微批风暴拓扑
【发布时间】:2018-12-17 18:13:43
【问题描述】:

如果我的问题重复,首先真诚道歉,我尝试搜索但找不到与我的问题相关的答案

首先真诚的道歉,如果我问一些非常基本的问题,因为我是 Storm 的初学者。 如果我的问题是重复的,当我尝试搜索但找不到与我的问题相关的答案时

请就我的以下用例提出建议。

我的用例:

我有一个 Spout 从一种内部消息传递机制读取数据,因为它以非常高的频率(100 秒/秒)接收和发送元组。

现在除了数据之外,每个元组也有一个频率(int)(因为总共可以有 4-5 种频率)。

现在我需要设计一个 Bolt 来批量/池化所有元组,并且仅按频率定期发出,具有仅发出最新元组的功能(以防在下一批之前收到重复),因为我们有一个基于字符串的键在元组数据中识别重复项。

例如

  1. 因此,所有以 25 秒为频率的元组将被汇集在一起​​,并由 Bolt 每 25 秒发出一次(如果在 25 秒内收到重复的元组,则只会考虑最新的一个)。

  2. 类似于所有以 10 分钟为频率的元组将被汇集在一起​​,并由 Bolt 每隔 10 分钟发出一次(如果在 10 分钟内收到重复的元组,则只会考虑最新的一个)。

** 现在,由于我们可以有 4-5 种频率(例如 10 秒、25 秒、10 分钟、20 分钟等,这些都是配置的),并且每个元组都应该被组合成一个适当的批次,并且发射(如上例)。

仅供参考。对于 Bolt 分组,我使用了“fieldsGrouping”,如下配置。

*.fieldsGrouping("FILTERING_BOLT",new Fields(PUBLISHING_FREQUENCY));*

请提供帮助或建议,什么是我的用例的最佳方法,因为我想不出任何东西来处理并发元组的流动和管理 Storm 的内部并行性。

【问题讨论】:

    标签: apache-storm


    【解决方案1】:

    听起来你想要窗口螺栓https://storm.apache.org/releases/2.0.0-SNAPSHOT/Windowing.html。可能你想要一个翻滚窗口(即窗口间隔之间没有重叠)

    Windowing bolts 让您设置它们应该发出窗口的时间间隔(例如每 10 秒),然后 Bolt 将在调用您提供的执行方法之前缓冲前 10 秒接收到的所有元组。

    我认为你想要的结构类似于例如

    spout -> splitter -> 5 second window bolt
                      -> 10 second window bolt
    

    拆分器应该接收元组,检查频率场并将元组发送到右侧窗口螺栓。您可以通过为每种频率类型声明一个流来做到这一点。

    public void declareOutputFields(OutputFieldsDeclarer declarer) {
        declarer.declare("5-sec-stream", ...);
        declarer.declare("10-sec-stream", ...);
    }
    
    public void execute(Tuple input) {
        if (frequencyIsFive(input)) {
            collector.emit("5-sec-stream", new Values(input.getValues()))
        }
        //more cases here
    }
    

    然后当你声明你的拓扑时

    topologyBuilder.setBolt("splitter", new SplitterBolt())
         .shuffleGrouping("spout")
    
    topologyBuilder.setBolt("5-second-window", new YourWindowingBolt())
         .globalGrouping("splitter", "5-sec-stream")
    

    使所有 5 秒元组都转到 5 秒窗口螺栓。

    请参阅https://storm.apache.org/releases/2.0.0-SNAPSHOT/Concepts.html 了解更多信息,尤其是有关流和分组的部分。

    https://github.com/apache/storm/blob/master/examples/storm-starter/src/jvm/org/apache/storm/starter/SlidingWindowTopology.java 有一个简单的窗口拓扑示例。

    您可能需要注意的一件事是 Storm 的元组超时。如果您需要一个窗口,例如10 分钟,您需要将元组超时从默认的 30 秒大幅提高,因此元组在队列中等待时不会超时。您可以通过设置例如conf.setMessageTimeoutSecs(15*60) 配置拓扑时。您希望在窗口间隔和元组超时之间有一点余地,因为您希望尽可能避免元组超时。

    【讨论】:

      猜你喜欢
      • 2016-08-21
      • 1970-01-01
      • 2013-08-06
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多