Apache Flink — Running Streams on a Real Engine
Slice the Same Orders with Four Window TVFs
Goal
Apply TUMBLE, HOP, CUMULATE, and SESSION windows to the same order stream, and confirm from the results how the window columns are attached and how many windows a row goes into. Build a window Top-N on top of a window aggregation and a window Top-N directly on top of a window TVF.
Why it matters
If you choose the wrong kind of window, the numbers are silently wrong — if you add up the counts of HOP, it is inflated by the amount of overlap, and if you leave the window columns out of GROUP BY, it is not a window aggregation but an unbounded aggregation and update logs get mixed in. A window TVF is only a function that attaches window columns to rows, so if you know exactly how those columns are attached, you can put either an aggregation or a ranking on top of it with ordinary SQL. The grader of this lab does not ask the cluster anything — it reads the sql-client output you saved and compares it with the expected values for each window that it computes directly from the source CSV in Python.
Steps
- Start the cluster with
flink-up, run /root/flink/windows/assign.sql, which picks only the orders withorder_id <= 20along with theorder_id, ts, window_start, window_end, window_timeattached by a 10-minuteTUMBLE, and save the output to /root/flink/windows/assign.out. - Save the output of /root/flink/windows/tumble.sql, which produces
cnt(count) andrevenue(sum of amount) per 10-minuteTUMBLEwindow and shop (shop), to /root/flink/windows/tumble.out. - Save the output of /root/flink/windows/hop.sql, which produces
cntfor eachHOP(slide 5 minutes, size 10 minutes) window, to /root/flink/windows/hop.out. - Save the output of /root/flink/windows/cumulate.sql, which produces
revenuefor eachCUMULATE(step 10 minutes, size 1 hour) window, to /root/flink/windows/cumulate.out. - Save the output of /root/flink/windows/session.sql, which produces
cntperSESSION(PARTITION BY shop, gap 5 minutes) window and shop, to /root/flink/windows/session.out. - Save the output of /root/flink/windows/top-shops.sql, which produces the top 2 shops by sales for each 10-minute TUMBLE window (
window_start, window_end, shop, revenue, rownum), to /root/flink/windows/top-shops.out. - Save the output of /root/flink/windows/top-orders.sql, which produces, without aggregation, the 3 orders with the largest amounts for each 30-minute TUMBLE window (
order_id, shop, amount, window_start, window_end, rownum), to /root/flink/windows/top-orders.out. - In /root/flink/windows/report.json, write
orders,hop_assignments,cumulate_windows,sessions, andmax_session_orders.
Notes
- Source:
/opt/lab/fixtures/data/windows_orders.csv, columnsorder_id BIGINT, shop STRING, amount INT, ts TIMESTAMP(3)(a CSV with no header, ts ascending). PutWATERMARK FOR ts AS ts - INTERVAL '1' SECONDon the table and run in streaming mode. - Shape:
FROM TUMBLE(TABLE orders, DESCRIPTOR(ts), INTERVAL '10' MINUTE). HOP is(TABLE, DESCRIPTOR, slide, size), CUMULATE is(TABLE, DESCRIPTOR, step, size), and SESSION is(TABLE orders PARTITION BY shop, DESCRIPTOR(ts), gap). - A window aggregation groups with
GROUP BY window_start, window_end, .... If you leave out the window columns, it becomes an unbounded aggregation and -U/+U get mixed in. - A common mistake: swapping the argument order of HOP and CUMULATE (the smaller value first). It is rejected with an error that size is not an integer multiple of slide (step).
- A common mistake: leaving window_start and window_end out of the PARTITION BY of a window Top-N — then it is not a window Top-N but a general Top-N, and a log comes out every time the ranking changes.
- Official docs: Windowing TVF · Window Aggregation · Window Top-N · Time Attributes
The three columns a window TVF attaches
Start the cluster with flink-up, run the source table orders (the columns and watermark from the Notes) and SELECT order_id, ts, window_start, window_end, window_time FROM TUMBLE(TABLE orders, DESCRIPTOR(ts), INTERVAL '10' MINUTE) WHERE order_id <= 20 in streaming in /root/flink/windows/assign.sql, and save the output to /root/flink/windows/assign.out.
A window TVF leaves the original columns as they are and returns them with three window columns attached. A window is a half-open interval [start, end), so an order stamped exactly at a boundary time goes to the window that starts at that time. Look at the difference between window_time and window_end.
A 10-minute aggregation per shop
In /root/flink/windows/tumble.sql, write a window aggregation that groups by a 10-minute TUMBLE window and shop and produces window_start, window_end, shop, COUNT(*) AS cnt, SUM(amount) AS revenue, and save the output to /root/flink/windows/tumble.out.
A window aggregation puts window_start and window_end in GROUP BY. The windows do not overlap, so adding up all the cnt values gives the original number of orders. The result comes out only once, as +I, when the window closes.
HOP — one order goes into two windows
In /root/flink/windows/hop.sql, write an aggregation that produces window_start, window_end, COUNT(*) AS cnt for each HOP(TABLE orders, DESCRIPTOR(ts), INTERVAL '5' MINUTE, INTERVAL '10' MINUTE) window, and save the output to /root/flink/windows/hop.out.
The third argument of HOP is slide (the interval at which windows start) and the fourth is size (the window length). With 10-minute windows starting every 5 minutes, the windows overlap by half and one order goes into two windows. Compare the sum of cnt with the original number of orders. The first window may start 5 minutes earlier than the first order.
CUMULATE — a cumulative window with a fixed start
In /root/flink/windows/cumulate.sql, write an aggregation that produces window_start, window_end, SUM(amount) AS revenue for each CUMULATE(TABLE orders, DESCRIPTOR(ts), INTERVAL '10' MINUTE, INTERVAL '1' HOUR) window, and save the output to /root/flink/windows/cumulate.out.
CUMULATE is a TUMBLE by size (1 hour) with the inside split into windows whose ends grow by step (10 minutes). The revenue of windows with the same start must grow or stay the same the later the end is. Count how many windows come out per hour.
SESSION — windows with a different length per shop
In /root/flink/windows/session.sql, write an aggregation that produces window_start, window_end, shop, COUNT(*) AS cnt per SESSION(TABLE orders PARTITION BY shop, DESCRIPTOR(ts), INTERVAL '5' MINUTE) window and shop, and save the output to /root/flink/windows/session.out.
A session continues if the gap between neighboring orders of the same shop is 5 minutes or less, and a new session starts if it is exceeded. The start of a session is the time of the first order, and the end is the last order + 5 minutes. If you leave out PARTITION BY, sessions are split with the shops mixed together.
A window Top-N on top of a window aggregation
In /root/flink/windows/top-shops.sql, compute the sales per 10-minute TUMBLE window and shop (SUM(amount) AS revenue), keep only the top 2 per window with ROW_NUMBER() OVER (PARTITION BY window_start, window_end ORDER BY revenue DESC) AS rownum, and produce window_start, window_end, shop, revenue, rownum. Save the output to /root/flink/windows/top-shops.out.
Put the window aggregation as a subquery, assign ROW_NUMBER on top of it, and filter with rownum <= 2 on the outside. Only when PARTITION BY has the two window columns does it become a window Top-N that emits the result just once when the window closes. A window with only one shop produces just one line.
A window Top-N directly on top of a window TVF
In /root/flink/windows/top-orders.sql, without aggregation, rank directly on top of the 30-minute TUMBLE window TVF with ROW_NUMBER() OVER (PARTITION BY window_start, window_end ORDER BY amount DESC) and produce the order_id, shop, amount, window_start, window_end, rownum of the 3 orders with the largest amounts for each window. Save the output to /root/flink/windows/top-orders.out.
A window Top-N can be put directly on top of the result of a window TVF, even without a window aggregation. Here the subject of the ranking is the order rows themselves, not aggregated rows. Do not use GROUP BY. The data was made with no ties in the amounts, so the ranking does not wobble.
Report — how many times is each counted per window
In /root/flink/windows/report.json, write as integers orders (the sum of cnt in tumble.out), hop_assignments (the sum of cnt in hop.out), cumulate_windows (the number of lines in cumulate.out), sessions (the number of lines in session.out), and max_session_orders (the maximum cnt in session.out).
Result rows start with '| +I |'. If you split with awk -F'|', column 1 is empty, column 2 is the op, and after that the SELECT columns follow in order. Compare hop_assignments with orders — it is a multiple of size/slide.