QuestionQ34

Ingesting and processing the data

You are analyzing a company’s stock price. Every 5 seconds, you must calculate a moving average for the prior 30 seconds of data. You read data from Pub/Sub and use Dataflow to perform the analysis. How should you configure the windowed pipeline?

Explanation

A 30-second sliding window that advances every 5 seconds creates overlapping 30-second intervals, yielding a moving average at five-second intervals. AfterWatermark.pastEndOfWindow() fires when the watermark passes the window’s end, producing the result for each event-time window. Apache Beam documents sliding windows as overlapping fixed-width time intervals and documents AfterWatermark.pastEndOfWindow() as firing when the watermark passes the end of the window.

Learn more

Community Discussion

No comments yet. Be the first to start the discussion!