Professional Data EngineerBuilding and operationalizing data processing systemsMedium
A data engineering team is building a pipeline that processes sensor data from IoT devices. The data arrives continuously and needs to be aggregated into 5-minute windows before being stored in BigQuery. Due to network fluctuations, some data points might arrive a few minutes late. The team needs to ensure that all data points, including late ones, are included in the correct 5-minute window without indefinitely delaying the pipeline. Which Apache Beam concept should they use to achieve this?
- ASliding Windows
- BFixed Windows
- CSession Windows
- DWatermarks and Allowed Lateness
Show answer & explanationAnswer & explanation
Correct answer: D. Watermarks and Allowed Lateness
Watermarks track the progress of event time, while allowed lateness specifies how long the system should wait for late data to arrive before finalizing a window. This combination ensures late data is included without indefinite delays.
Why the other options are wrong
- A. Sliding windows overlap and are useful for calculating aggregates over moving time periods, but don't solve the late data problem directly.
- B. Fixed windows define contiguous, non-overlapping windows but don't inherently handle late data gracefully.
- C. Session windows group data based on a gap of inactivity, which is not suitable for fixed-duration aggregations with late data.
Watermarks and Allowed Lateness
In Apache Beam, watermarks are a measure of processing progress based on event time, and allowed lateness defines a grace period after a window's watermark passes during which late data can still be processed for that window.
- Watermarks indicate event-time completeness
- Allowed lateness handles data arriving after its window's end
- Prevents indefinite delays while ensuring data accuracy
Memory trick: Watermarks and Lateness: The clock and the grace period.