【发布时间】:2021-05-03 03:23:38
【问题描述】:
我的事件是这样的:case class Event(user: User, stats: Map[StatType, Int])
每个事件都包含 +1 或 -1 值。 我目前的管道运行良好,但每次统计数据更改都会产生新事件。
eventsStream
.keyBy(extractKey)
.reduce(reduceFunc)
.map(prepareRequest)
.addSink(sink)
我想在一个时间窗口中聚合这些增量,然后再将它们与当前状态合并。所以我想要同样的滚动减少,但有一个时间窗口。
当前简单滚动减少:
500 – last reduced value
+1
-1
+1
Emitted events: 501, 500, 501
带窗口的滚动减少:
500 – last reduced value
v-- window
+1
-1
+1
^-- window
Emitted events: 501
我尝试过简单的解决方案,将时间窗口放在 reduce 之前,但在阅读文档后,我发现 reduce 现在有不同的行为。
eventsStream
.keyBy(extractKey)
.timeWindow(Time.minutes(2))
.reduce(reduceFunc)
.map(prepareRequest)
.addSink(sink)
看来我应该在减少时间窗口后制作键控流并减少它:
eventsStream
.keyBy(extractKey)
.timeWindow(Time.minutes(2))
.reduce(reduceFunc)
.keyBy(extractKey)
.reduce(reduceFunc)
.map(prepareRequest)
.addSink(sink)
它是解决问题的正确管道吗?
【问题讨论】:
-
实际上,将窗口放在
reduce之前时是否有任何问题或错误消息? AFAIK 应该可以工作。 -
在流中我有像
case class Event(user: User, stats: Map[StatType, Int])这样的事件。每个事件都包含 +1 或 -1 值。正如我在文档reduce中读到的那样,键控流会发出一个新状态。因此,如果我为某些用户和统计类型设置了 500 的值,如果流中有 +1 事件,它将发出 501。但是应用于窗口流的 reduce 仅减少窗口内的那些事件。所以看起来它会发出增量而不是新状态。
标签: scala apache-flink flink-streaming