默认情况下,Apache Beam会在无界PCollection中为每个元素创建一个全局窗口。这意味着,每个元素都属于唯一的窗口,且窗口的边界为无限大。由于全局窗口没有边界,因此Beam Runner将在处理数据时不断等待新的数据到达。这种行为称为无限等待。
以下是一个基于Java的示例代码,展示了如何创建一个无界PCollection以及默认行为的实现:
Pipeline pipeline = Pipeline.create();
PCollection unboundedCollection =
pipeline.apply(TextIO.read().from("input.txt"))
.apply(Window.into(new GlobalWindows())
.triggering(AfterWatermark.pastEndOfWindow())
.withAllowedLateness(Duration.standardMinutes(10))
.discardingFiredPanes())
.apply(MapElements.via(new SimpleFunction() {
@Override
public String apply(String input) {
return input.toUpperCase();
}
}));
unboundedCollection.apply(TextIO.write().to("output.txt"));
在上述示例代码中,我们通过应用 GlobalWindows()
方法来创建了一个无界PCollection的全局窗口。同时,我们也可以看到,定义了一个默认的触发器:AfterWatermark.pastEndOfWindow()
。这意味着,我们想要在收到任何新元素时处理窗口,而无需考虑是否已到达窗口边界。默认的 allowed lateness 为 0,firing mode 为 discarding,表示元素已被丢弃。
注意,在没有显式定义触发器的情况下,Beam Runner将使用默认触发器:Repeatedly.forever(AfterPane.elementCountAtLeast(1))
,这意味着在每个窗格中累积任意数量的元素时,Beam Runner将触发窗口处理。这对于有限数据集非常有用,但对于无限数据集来说则可能会导致问题,因为Beam Runner将不会在新数据到达时处理窗口。
因此,如果我们希望程序在一段时间内处理窗口并输出结果,我们需要显式地定义触发器,并将 allowed lateness 设置为适当