In stream processing with systems like Kafka and Flink, effectively managing late data, understanding ordering guarantees, and handling backpressure are critical for building robust and accurate pipelines. Late data refers to events that arrive at the processing system after their logical 'event time' window has supposedly closed. This often occurs due to network latency, clock skews, or upstream processing delays. For example, a transaction event from a mobile device might arrive minutes after its timestamp due to poor connectivity. Flink addresses this primarily through watermarks combined with an 'allowed lateness' threshold, which enables windows to temporarily hold open for a specified duration, ensuring late but still relevant events are included in aggregations, preventing data incompleteness at the cost of slight result latency.
Key Takeaways
- Late data requires explicit strategies (e.g., Flink's watermarks with allowed lateness) to ensure correct windowed aggregations and data completeness.
- Ordering guarantees in distributed stream processing are typically per-partition or per-key, not global, allowing for scalability.
- Backpressure is an essential self-preservation mechanism that prevents system overload by signaling upstream components to slow down.
- Monitoring and understanding backpressure is crucial for identifying and addressing bottlenecks in your streaming pipeline.
Code Example
import org.apache.flink.api.common.eventtime.WatermarkStrategy;
import org.apache.flink.streaming.api.datastream.DataStream;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import org.apache.flink.streaming.api.windowing.assigners.TumblingEventTimeWindows;
import org.apache.flink.streaming.api.windowing.time.Duration;
public class LateDataExample {
public static void main(String[] args) throws Exception {
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
DataStream<SensorReading> readings = env.fromElements(
new SensorReading("sensor_1", 1678886400000L, 25.0), // 2023-03-15 00:00:00
new SensorReading("sensor_1", 1678886405000L, 26.0), // 2023-03-15 00:05:00
new SensorReading("sensor_1", 1678886462000L, 27.0) // 2023-03-15 00:01:02 - LATE for 00:00 window
)
.assignTimestampsAndWatermarks(WatermarkStrategy
.<SensorReading>forBoundedOutOfOrderness(Duration.ofSeconds(5)) // Allow events up to 5 seconds late
.withTimestampAssigner((event, timestamp) -> event.timestamp));
readings.keyBy(r -> r.id)
.window(TumblingEventTimeWindows.of(Duration.ofMinutes(1))) // 1-minute windows
.sum("value")
.print();
env.execute("Late Data Example");
}
// SensorReading class definition omitted for brevity, assumes Long timestamp and String id
}How this code works
This Flink program illustrates how stream processing systems manage "late" data—events that arrive out of order relative to their timestamps. It sets up a DataStream using a few SensorReading events, specifically including one with a timestamp of 00:01:02 that arrives after an event timestamped 00:05:00, creating an out-of-order scenario that tests the system's ability to correctly process events despite their arrival sequence.
The core of the solution lies in assignTimestampsAndWatermarks. Here, WatermarkStrategy.forBoundedOutOfOrderness(Duration.ofSeconds(5)) is used. This strategy configures Flink to generate "watermarks"—markers of event time progress—allowing events up to 5 seconds "late" to still be processed within their respective windows. The data is then grouped by keyBy(r -> r.id) and processed in TumblingEventTimeWindows.of(Duration.ofMinutes(1)), meaning one-minute non-overlapping windows, before sum("value") aggregates the readings. A subtle point for beginners: while a grace period is defined, in this specific example, the 00:05:00 event arriving early will cause the watermark to advance significantly. As a result, the 00:01:02 event might still be considered too late even with the 5-second allowance, and could be dropped, highlighting the importance of appropriately sizing the forBoundedOutOfOrderness duration for actual data patterns.