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?

  1. ASliding Windows
  2. BFixed Windows
  3. CSession Windows
  4. DWatermarks and Allowed Lateness
Show answer & 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.

More Building and operationalizing data processing systems questions