TT Lab
Get started
Learn Learning paths Courses

Apache Flink — Running Streams on a Real Engine

Slice the Same Orders with Four Window TVFs

Continue in TT Lab

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

  1. Start the cluster with flink-up, run /root/flink/windows/assign.sql, which picks only the orders with order_id <= 20 along with the order_id, ts, window_start, window_end, window_time attached by a 10-minute TUMBLE, and save the output to /root/flink/windows/assign.out.
  2. Save the output of /root/flink/windows/tumble.sql, which produces cnt (count) and revenue (sum of amount) per 10-minute TUMBLE window and shop (shop), to /root/flink/windows/tumble.out.
  3. Save the output of /root/flink/windows/hop.sql, which produces cnt for each HOP (slide 5 minutes, size 10 minutes) window, to /root/flink/windows/hop.out.
  4. Save the output of /root/flink/windows/cumulate.sql, which produces revenue for each CUMULATE (step 10 minutes, size 1 hour) window, to /root/flink/windows/cumulate.out.
  5. Save the output of /root/flink/windows/session.sql, which produces cnt per SESSION (PARTITION BY shop, gap 5 minutes) window and shop, to /root/flink/windows/session.out.
  6. 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.
  7. 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.
  8. In /root/flink/windows/report.json, write orders, hop_assignments, cumulate_windows, sessions, and max_session_orders.

Notes

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.