【发布时间】:2018-11-07 18:15:49
【问题描述】:
我正在使用 Chicago Traffic Tracker 数据集,其中每 15 分钟发布一次新数据。当有新数据可用时,它代表距离“实时”(example,查找_last_updt)10-15 分钟的记录。
例如,在 00:20,我得到时间戳为 00:10 的数据;在 00:35,我从 00:20 开始;在 00:50,我从 00:40 开始。因此,我可以“固定”获取新数据的间隔(每 15 分钟一次),尽管时间戳的间隔略有变化。
我正在尝试在 Dataflow (Apache Beam) 上使用这些数据,为此我正在使用滑动窗口。我的想法是收集和处理 4 个连续的数据点(4 x 15 分钟 = 60 分钟),理想情况下,一旦有新的数据点可用,就更新我对总和/平均值的计算。为此,我从代码开始:
PCollection<TrafficData> trafficData = input
.apply("MapIntoSlidingWindows", Window.<TrafficData>into(
SlidingWindows.of(Duration.standardMinutes(60)) // (4x15)
.every(Duration.standardMinutes(15))) . // interval to get new data
.triggering(AfterWatermark
.pastEndOfWindow()
.withEarlyFirings(AfterProcessingTime.pastFirstElementInPane()))
.withAllowedLateness(Duration.ZERO)
.accumulatingFiredPanes());
不幸的是,当我从我的输入中接收到一个新的数据点时,我没有从我所追求的GroupByKey 获得新的(更新的)结果。
我的 SlidingWindows 有问题吗?还是我错过了什么?
【问题讨论】:
-
您的意思是在第一个元素之后没有得到任何元素,还是在第一次触发后没有得到添加到窗口的后期元素?如果是后者,那么很可能是
allowedLateness(Duration.ZERO)引起的,这会丢弃所有的后期元素。 -
嗨@Anton,我在第一次发射后没有得到迟到的元素,即使这些元素应该在同一个“窗口”上。例如,在 01:14 到达的元素应该包含在从 00:15 开始的窗口中,但事实并非如此。我对
allowedLateness的理解是,将其设置为大于 0(比如说 5 分钟),将允许包含在预计关闭窗口之后到达的元素(因此,如果 01:14 的元素恰好在 01 到达:18,它仍将包含在 01:15 关闭的窗口中)。如果我的理解有误,请告诉我。
标签: java google-cloud-dataflow apache-beam sliding-window