A data engineer needs to design a stream processing pipeline that reads events from Pub/Sub, enriches them with data from a Cloud Storage file, and writes aggregated results to BigQuery. The pipeline must handle late-arriving events up to 1 hour. Which Dataflow feature should be used to manage late data?
Trap 1: Triggers
Triggers control when aggregate results are emitted, but watermarks signal late data.
Trap 2: Side inputs
Side inputs are for enriching streams with static or slowly changing data, not for managing late data.
Trap 3: Windowing
Windowing groups elements by time, but watermarks handle late arrivals.
- A
Triggers
Why it fails: Triggers control when aggregate results are emitted, but watermarks signal late data.
- B
Watermarks
Watermarks track event-time progress in Dataflow and, combined with allowed lateness, let the pipeline retain windows and process events arriving up to one hour late. This directly satisfies the requirement to handle late-arriving Pub/Sub events.
- C
Side inputs
Why it fails: Side inputs are for enriching streams with static or slowly changing data, not for managing late data.
- D
Windowing
Why it fails: Windowing groups elements by time, but watermarks handle late arrivals.