【发布时间】:2018-12-17 18:13:43
【问题描述】:
如果我的问题重复,首先真诚道歉,我尝试搜索但找不到与我的问题相关的答案
首先真诚的道歉,如果我问一些非常基本的问题,因为我是 Storm 的初学者。 如果我的问题是重复的,当我尝试搜索但找不到与我的问题相关的答案时
请就我的以下用例提出建议。
我的用例:
我有一个 Spout 从一种内部消息传递机制读取数据,因为它以非常高的频率(100 秒/秒)接收和发送元组。
现在除了数据之外,每个元组也有一个频率(int)(因为总共可以有 4-5 种频率)。
现在我需要设计一个 Bolt 来批量/池化所有元组,并且仅按频率定期发出,具有仅发出最新元组的功能(以防在下一批之前收到重复),因为我们有一个基于字符串的键在元组数据中识别重复项。
例如
-
因此,所有以 25 秒为频率的元组将被汇集在一起,并由 Bolt 每 25 秒发出一次(如果在 25 秒内收到重复的元组,则只会考虑最新的一个)。
类似于所有以 10 分钟为频率的元组将被汇集在一起,并由 Bolt 每隔 10 分钟发出一次(如果在 10 分钟内收到重复的元组,则只会考虑最新的一个)。
** 现在,由于我们可以有 4-5 种频率(例如 10 秒、25 秒、10 分钟、20 分钟等,这些都是配置的),并且每个元组都应该被组合成一个适当的批次,并且发射(如上例)。
仅供参考。对于 Bolt 分组,我使用了“fieldsGrouping”,如下配置。
*.fieldsGrouping("FILTERING_BOLT",new Fields(PUBLISHING_FREQUENCY));*
请提供帮助或建议,什么是我的用例的最佳方法,因为我想不出任何东西来处理并发元组的流动和管理 Storm 的内部并行性。
【问题讨论】:
标签: apache-storm