Schiebefenster werden nativ in Beam unterstützt. Bitte beachten Sie die programming guide und Dokumentation für die SlidingWindows Klasse.
Z. B .:
PCollection<Foo> foos = ...;
PCollection<Integer> counts = foos
.apply(Window.into(
SlidingWindows.of(Duration.standardMinutes(5))
.every(Duration.standardMinutes(1))))
// Below is required instead of Count.globally() when you use
// a non-global windowing function.
.apply(Combine.globally(Count.<Foo>combineFn()).withoutDefaults());
PCollection<String> formattedCounts = counts.apply(
ParDo.of(new DoFn<Integer, String>() {
@ProcessElement
public void process(ProcessContext c, BoundedWindow w) {
c.output("Window: " + w + ", count: " + c.element());
}
}));
Auslösung ist eine separate Dimension des Problems, und sie steuert, wenn die Daten für ein bestimmtes Fenster wird „vollständig genug“ betrachtet werden, die Aggregation anzuwenden. Siehe programming guide.