【发布时间】:2017-05-04 16:32:33
【问题描述】:
我想知道在 Apache Storm 中解决以下问题的最佳做法是什么。
我有一个单独的 spout,它生成一个带有明确时间戳的整数值流。目标是在此流上使用三个滑动窗口执行最小/最大聚合:
- 最后一小时
- 最后一天,即过去 24 小时
最后一小时很容易:
topology.setBolt("1h", ...)
.shuffleGrouping("spout")
.withWindow(Duration.hours(1), Duration.seconds(10))
.withTimestampField("timestamp"));
但是,在较长时期内,我担心窗口的队列大小。当我像上一小时聚合一样直接从 spout 使用元组时,每个元组都会最终进入队列。
一种可能性是从预先聚合的“1h”bolt 中消耗元组。但是,由于我使用的是显式时间戳,因此来自“1h”螺栓的迟到的元组将被忽略。 1 小时的延迟不是一个选项,因为这会延迟窗口的评估。有没有办法“允许”迟到的元组而不影响结果的及时性?
当然,我也可以每小时存储一个聚合,然后计算过去 24 小时的最小值,包括“1h”流中的最新值。但我很好奇是否有办法使用 Storm 方法正确地做到这一点。
更新 1
感谢 arunmahadevan 的回答,我更改了 1h min 螺栓以在相应的 1h 窗口中发出所有元组的最大时间戳的最小元组。这样,消费螺栓不会因为迟到而丢弃元组。我还引入了一个新字段original-timestamp 来保留最小元组的原始时间戳。
更新 2
我终于找到了一种更好的方法,即只在 1 小时分钟的螺栓中发出状态变化。只要没有收到新的元组,Storm 就不会在消耗螺栓中提前时间,因此可以防止迟到问题。此外,我可以保留原始时间戳,而无需将其复制到单独的字段中。
【问题讨论】:
标签: apache-storm