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?
A Use a fixed window with a duration of 5 seconds. Emit results by setting the following trigger: AfterProcessingTime.pastFirstElementInPane().plusDelayOf (Duration.standardSeconds(30)) B Use a fixed window with a duration of 30 seconds. Emit results by setting the following trigger: AfterWatermark.pastEndOfWindow().plusDelayOf (Duration.standardSeconds(5)) C Use a sliding window with a duration of 5 seconds. Emit results by setting the following trigger: AfterProcessingTime.pastFirstElementInPane().plusDelayOf (Duration.standardSeconds(30)) D Use a sliding window with a duration of 30 seconds and a period of 5 seconds. Emit results by setting the following trigger: AfterWatermark.pastEndOfWindow () Show Answer Answer 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