A sliding time window also represents time intervals in the data stream; however, sliding time windows can overlap. For example, each window might capture 60 seconds' worth of data, but a new window starts every 30 seconds. The frequency with which sliding windows begin is called the period. Therefore, our example would have a window duration of 60 seconds and a period of 30 seconds.
Because multiple windows overlap, most elements in a data set will belong to more than one window. This kind of windowing is helpful for taking running data averages; using sliding time windows, you can compute a running average of the past 60 seconds’ worth of data, updated every 30 seconds.
The following example code shows how to apply Window to divide a PCollection into sliding time windows. Each window is 30 seconds in length, and a new window begins every five seconds:
{{if (eq .Sdk “go”)}}
slidingWindowedItems := beam.WindowInto(s, window.NewSlidingWindows(5*time.Second, 30*time.Second), input)
{{end}}
{{if (eq .Sdk “java”)}}
PCollection<String> input = ...;
PCollection<String> slidingWindowedItems = input.apply(
Window.<String>into(SlidingWindows.of(Duration.standardSeconds(30)).every(Duration.standardSeconds(5))));
{{end}} {{if (eq .Sdk “python”)}}
from apache_beam import window
sliding_windowed_items = (
input | 'window' >> beam.WindowInto(window.SlidingWindows(30, 5)))
{{end}}
Because multiple windows overlap, most elements in a data set will belong to more than one window. This kind of windowing is useful for taking running averages of data; using sliding time windows, you can compute a running average of the past 60 seconds’ worth of data, updated every 30 seconds, in our example.
{{if (eq .Sdk “go”)}}
max := stats.Max(s, windowedData) mean := stats.Mean(s, windowedData) min := stats.Min(s, windowedData)
{{end}} {{if (eq .Sdk “java”)}}
Combine.globally(Max.ofIntegers()) Combine.globally(Mean.ofIntegers()) Combine.globally(Min.ofIntegers())
{{end}} {{if (eq .Sdk “python”)}}
beam.CombineGlobally.globally(beam.combiners.MaxCombineFn()) beam.CombineGlobally.globally(beam.combiners.MeanCombineFn()) beam.CombineGlobally.globally(beam.combiners.MinCombineFn())
{{end}}