Professional Data EngineerBuilding and operationalizing data processing systemsHard
A data engineering team is troubleshooting a Dataflow job that processes real-time sensor data. They observe that some incoming data records are processed with significant delays, leading to outdated results in their analytics dashboard. The Dataflow monitoring interface shows a consistently increasing 'System Lag' metric, while 'Data freshness' is also high. The pipeline uses fixed windows and processes data based on event time. What is the most likely cause of this issue?
- AThe BigQuery sink is experiencing write bottlenecks.
- BThe Pub/Sub source is not delivering messages fast enough.
- CHigh watermark lag due to out-of-order or late-arriving data.
- DInsufficient Dataflow worker CPU or memory.
Show answer & explanationAnswer & explanation
Correct answer: C. High watermark lag due to out-of-order or late-arriving data.
Consistently increasing 'System Lag' and high 'Data freshness' (meaning data is old) in an event-time-based streaming pipeline strongly indicates that the watermark is not advancing quickly enough, typically due to a high volume of out-of-order or late-arriving data, causing the system to wait for expected data.
Why the other options are wrong
- A. A BigQuery write bottleneck would show up as high sink latency or errors, not primarily as increasing 'System Lag' related to event time.
- B. A slow Pub/Sub source would reduce overall throughput but wouldn't typically cause an *increasing* 'System Lag' related to event time, unless it consistently delivered extremely late data.
- D. Insufficient worker resources would typically manifest as high CPU utilization, high memory usage, or 'Stuck' workers, impacting processing time more broadly than specifically event-time lag.
Dataflow Watermark Lag
Watermark lag in Dataflow indicates the difference between the current processing time and the event time of the latest data processed, signifying how far behind the system is in processing event-time data.
- High lag often caused by late-arriving or out-of-order data.
- Affects data freshness and accuracy of event-time windowed computations.
- Monitored via 'System Lag' and 'Data freshness' metrics.
- Can be mitigated by `withAllowedLateness` and proper windowing strategies.
Memory trick: Lagging system, old data, watermark's slow, problems flow.