Apache Flink — Running Streams on a Real Engine
Count the Dropped Rows While Changing the Delay
Goal
Run the same data in batch and in streaming with several delays, and confirm with numbers when an event-time window closes according to the watermark and which rows it drops. Find the late rows yourself with CURRENT_WATERMARK and compare them with the rows the window drops.
Why it matters
A row that arrives after the watermark has closed the window disappears without an error. So a problem like "the numbers are a bit short" cannot be found in logs and can be explained only if you know how the watermark moves. The input of this lab has second-granularity timestamps, so the watermark advances immediately at every record, and the result is determined by the arrival order alone, regardless of execution speed. The grader does not ask the cluster anything — it reads the sql-client output you saved, feeds the source CSV through in arrival order, and compares with the values it computes using the same rules as the engine (emit the row first and then advance the watermark · drop the rows of a window with window_end ≤ 워터마크, that is, a window end at or below the watermark).
Steps
- Start the cluster with
flink-up, write aneventstable withWATERMARK FOR ts AS ts - INTERVAL '5' SECONDandDESCRIBE events;in /root/flink/watermark/ddl.sql, and save the output to /root/flink/watermark/ddl.out. - Run /root/flink/watermark/batch.sql, which in batch mode produces
cnt(count) andtotal(sum of reading) for each 1-minuteTUMBLEwindow, and save the output to /root/flink/watermark/batch.out. - Save the output of /root/flink/watermark/w5.sql, which runs the same aggregation in streaming mode (5-second delay), to /root/flink/watermark/w5.out.
- In /root/flink/watermark/sweep.sql, create a table
events_0whose watermark istsand a tableevents_30whose watermark ists - INTERVAL '30' SECOND, run the same aggregation in the orderevents_0→events_30, and save the output to /root/flink/watermark/sweep.out. - Save the output of /root/flink/watermark/late.sql, which picks the
event_id, ts, wmof rows whereCURRENT_WATERMARK(ts)is not NULL andts <= CURRENT_WATERMARK(ts)from the 5-second delayevents, to /root/flink/watermark/late.out. - Save the output of /root/flink/watermark/filtered.sql, which produces the same 1-minute aggregation after filtering out the late rows in front of the window (
CURRENT_WATERMARK(ts) IS NULL OR ts > CURRENT_WATERMARK(ts)), to /root/flink/watermark/filtered.out. - Save the output of /root/flink/watermark/slow.sql, which is the same aggregation as step 3 with only
SET 'pipeline.auto-watermark-interval' = '1 h';added, to /root/flink/watermark/slow.out. - In /root/flink/watermark/report.json, write
total_rows,dropped_0,dropped_5,dropped_30,late_rows_5, anddropped_slow.
Notes
- Source columns:
event_id BIGINT, sensor STRING, reading INT, ts TIMESTAMP(3)(a CSV with no header, file/opt/lab/fixtures/data/watermark_events.csv, file order = arrival order). - Shape of the window aggregation:
SELECT window_start, window_end, COUNT(*) AS cnt, SUM(reading) AS total FROM TUMBLE(TABLE 표, DESCRIPTOR(ts), INTERVAL '1' MINUTE) GROUP BY window_start, window_end;(the placeholder stands for the table). - If one SQL file has several SELECTs, jobs run in turn, and the result tables are printed in order in the output.
- A common mistake: batch mode does not use watermarks. To see late rows, you have to run in streaming. Leave the parallelism at the default 1 — if there are several inputs, the watermark follows the minimum among them.
- A common mistake: at the first row, there is no watermark yet, so
CURRENT_WATERMARK(ts)is NULL. A comparison with NULL is not true and does not pass the WHERE, so if you leave out theIS NULLcondition from the filtering expression of step 6, even the first row is dropped. - Official docs: Timely Stream Processing · CREATE — WATERMARK · Time Attributes · Windowing TVF · Built-in Functions · Configuration
Declare ts as the event time
Start the cluster with flink-up, write an events table that reads the source CSV (columns event_id BIGINT, sensor STRING, reading INT, ts TIMESTAMP(3), WATERMARK FOR ts AS ts - INTERVAL '5' SECOND) and DESCRIBE events; in /root/flink/watermark/ddl.sql, and save the output to /root/flink/watermark/ddl.out.
Put the WATERMARK clause inside the column list, after the last column. If ROWTIME appears next to the type of ts in the DESCRIBE result and the expression appears in the watermark column, it has become an event-time attribute. The path of the filesystem connector is file:///opt/lab/fixtures/data/watermark_events.csv.
Build the baseline in batch
In /root/flink/watermark/batch.sql, write SET 'execution.runtime-mode' = 'batch';, the events table from step 1, and an aggregation that produces window_start, window_end, COUNT(*) AS cnt, SUM(reading) AS total for each 1-minute TUMBLE window, and save the output to /root/flink/watermark/batch.out.
Batch calculates after gathering all the input, so it drops no rows because of the watermark. That is why this result becomes the baseline for "if there had been no late rows at all". TUMBLE can be used on a TIMESTAMP column in batch too. The sum of all the cnt values must equal the number of rows in the source.
Streaming with a 5-second delay — the rows that get dropped
Create /root/flink/watermark/w5.sql, which runs the same aggregation as step 2 with SET 'execution.runtime-mode' = 'streaming'; (the watermark stays a 5-second delay), and save the output to /root/flink/watermark/w5.out.
A window emits its result once the moment its window_end is at or below the watermark, and clears its state. A row that would belong to that window and arrives after that is dropped. Compare the sum of cnt in the result table with batch. Window results come out only as +I.
Compare 0-second and 30-second delays at once
In /root/flink/watermark/sweep.sql, create a table events_0 whose watermark is ts and a table events_30 whose watermark is ts - INTERVAL '30' SECOND (columns and source the same as events), run the same 1-minute aggregation in streaming mode with events_0 first and events_30 next, and save the output to /root/flink/watermark/sweep.out.
A watermark expression is attached per table, so to change the delay you create separate tables. Two SELECTs in one file run as two jobs in turn, and the result tables are printed in that order in the output. With a delay of 0, even a slightly late row is dropped, and with 30 seconds it waits for most of them.
Pick out the late rows with CURRENT_WATERMARK
Create /root/flink/watermark/late.sql, which runs SELECT event_id, ts, CURRENT_WATERMARK(ts) AS wm ... WHERE CURRENT_WATERMARK(ts) IS NOT NULL AND ts <= CURRENT_WATERMARK(ts) in streaming mode on the 5-second delay events, and save the output to /root/flink/watermark/late.out.
CURRENT_WATERMARK is the current watermark of the operator the row passes through. The watermark generator emits the row first and then advances the watermark, so the value a row sees is the maximum ts up to the rows before it − 5 seconds. Compare the number of late rows with the number of rows dropped in step 3 — being late does not mean the window is closed.
Filtering in front of the window drops more
Create /root/flink/watermark/filtered.sql, which keeps only rows with CURRENT_WATERMARK(ts) IS NULL OR ts > CURRENT_WATERMARK(ts) from the 5-second delay events (a view or subquery) and then produces the same 1-minute TUMBLE aggregation in streaming, and save the output to /root/flink/watermark/filtered.out.
This is the expression the docs recommend for filtering out late rows. It looks row by row at "is it earlier than the watermark", so it also drops rows that the window would have accepted because it was still open. To put a view into the input of a window TVF, write it as TUMBLE(TABLE view, ...) after CREATE VIEW.
A watermark interval of 1 hour — the advance stops
Run /root/flink/watermark/slow.sql, which is w5.sql from step 3 with only the line SET 'pipeline.auto-watermark-interval' = '1 h'; added at the very front, and save the output to /root/flink/watermark/slow.out.
A watermark goes out at every period (this setting), and in between, it goes out right at a record only when the new value is ahead of the last emitted value by more than the interval. With an interval of 1 hour, after the first watermark (which goes out right away because nothing was emitted before it), it does not advance until the file ends. Then will there be any window that closes?
Report — the delay and the dropped rows
In /root/flink/watermark/report.json, write as integers total_rows (the sum of cnt in batch.out), dropped_0, dropped_5, and dropped_30 (the batch sum minus the sum of cnt of each delay), late_rows_5 (the number of rows in late.out), and dropped_slow (the batch sum − the sum of cnt in slow.out).
You can count all of them from the saved output. Result rows start with '| +I |', and batch result rows start with a date. Add up the cnt column with awk -F'|'. sweep.out has two result tables, so split and count them using the header line (the line with op) as the boundary.