【发布时间】:2018-02-20 13:57:15
【问题描述】:
我怎样才能让 GroupByKey 触发早期结果,而不是等待所有数据到达(在我的情况下这是一个相当长的时间)。我尝试将我的输入 PCollection 拆分为具有早期触发器的窗口,但它只是行不通。在给出结果之前,它仍然等待所有数据到达。
PCollection<List<String>> input = ...
PCollection<KV<Integer,List<String>>> keyedInput = input.apply(ParDo.of(new AddArbitraryKey()))
keyedInput.apply(Window<KV<Integer,List<String>>>into(
FixedWindows.of(Duration.standardSeconds(1)))
.triggering(Repeatedly.forever(AfterWatermark.pastEndOfWindow()))
.withAllowedLateness(Duration.ZERO).discardingFiredPanes())
.apply(GroupByKey.<Integer,List<String>>create())
.apply(ParDo.of(new RemoveArbitraryKey()))
.apply(ParDo.of(new FurtherProcessing())
我这样做是为了防止 fusing 。 AddArbitraryKey 转换使用时间戳输出其元素。但是,GroupByKey 会保留所有内容,直到所有数据到达(对于所有窗口)。有人可以告诉我如何让它尽早触发。谢谢你 。
【问题讨论】:
-
我刚刚遇到了同样的问题,一个简单的固定窗口没有帮助,一个组合每个键等到有界集合结束。你找到解决办法了吗?
标签: google-cloud-dataflow apache-beam