Apache Flink — Running Streams on a Real Engine
Joining orders three ways
Goal
Join users, deliveries, and exchange rates to the same order flow with a regular join, an interval join, and an event-time temporal join, and confirm from the output and the execution plans what each join remembers and whether it fixes its results.
Why it matters
A stream join piles up rows of the other side in state because it does not know when the partner will arrive. With no time condition it holds them forever, with a time range until the watermark passes, and with a versioned table only the needed versions. If you choose the wrong join type, state grows endlessly or past results silently change. The grader of this lab does not ask the cluster anything; it compares the sql-client output and plan JSON you saved with values it computes directly from the source CSV.
Steps
- Start the cluster with
flink-upand define the four sources (users · orders · shipments · rates) in /root/flink/joins/ddl.sql. Put watermarks on the time columns of orders, shipments, and rates. Save the output of /root/flink/joins/count.sql, which counts the rows of the four tables in batch, to /root/flink/joins/count.out (column namestblandn). - In /root/flink/joins/regular.sql, write a streaming query that joins orders and users on
user_idand producesorders(count) andamount(sum) by tier, and save the output to /root/flink/joins/regular.out. - In /root/flink/joins/interval.sql, write an interval join that attaches deliveries that went out within 2 hours of the order time (both ends inclusive), and save
order_id,ship_id, anddelay_s(seconds) to /root/flink/joins/interval.out. - In /root/flink/joins/unshipped.sql, write an outer interval join that produces the
order_idandorder_timeof orders that were not delivered within 2 hours, and save the output to /root/flink/joins/unshipped.out. - Make a versioned view
rates_vfrom rates, and save the output of /root/flink/joins/temporal.sql, a temporal join that attaches the exchange rate at the order time and producesorder_id,currency,rate, andamount_krw(amount × rate), to /root/flink/joins/temporal.out. - Save the output of /root/flink/joins/latest.sql, which joins the same
rates_vwithoutFOR SYSTEM_TIME AS OFand producesordersandtotal_krwper currency, to /root/flink/joins/latest.out. - Extract the execution plans of the three joins (steps 2, 3, and 5) with
COMPILE PLANto /root/flink/joins/regular-plan.json, /root/flink/joins/interval-plan.json, and /root/flink/joins/temporal-plan.json. - In /root/flink/joins/report.json, write
matched_orders,interval_rows,unshipped_orders,temporal_rows,orders_without_rate, andusd_gap_krw.
Notes
- Sources (CSVs with no header, times in seconds, each file sorted in its own time order):
joins_users.csv=user_id, tier, signup_time·joins_orders.csv=order_id, user_id, currency, amount, order_time·joins_shipments.csv=ship_id, order_id, ship_time·joins_rates.csv=currency, rate, update_time(rate is in won; read it asDECIMAL(10, 4)). All of them are in/opt/lab/fixtures/data/. - To write the definitions only once, use
sql-client.sh -i ddl.sql -f 쿼리.sql > 쿼리.out 2>&1(the placeholder stands for the query file name) — the CREATE statements in the-ifile are run first. - Streaming results have an
opcolumn (+I · -U · +U · -D) at the front. The job ends at the end of the file, and then the watermark advances all the way, so all the remaining intervals and versions are processed. COMPILE PLANtakes anINSERT INTOstatement. Create and use a sink with'connector' = 'blackhole'. If a file already exists at the same path, it raises an error instead of overwriting, so delete it first when extracting again.- A common mistake: if you put the time condition of an outer interval join in
WHERE, the null rows get filtered out. The right side of a temporal join needs a primary key, and an append-only source has to be made into one with a deduplication view. - Official docs: Joins · Versioned Tables · Deduplication · SQL Client
Define the four sources and count the rows
After flink-up, define users, orders, shipments, and rates in /root/flink/joins/ddl.sql (with watermarks on order_time, ship_time, and update_time). Run /root/flink/joins/count.sql, which produces the row counts of the four tables in batch mode as columns tbl and n, with sql-client.sh -i ddl.sql -f count.sql, and save it to /root/flink/joins/count.out.
A watermark is written like WATERMARK FOR time_column AS time_column. Each file is sorted in its own time order, so there are no late rows even without giving a delay. If you join the four SELECTs with UNION ALL, it comes out as one table.
Regular join — remember both sides without looking at time
In /root/flink/joins/regular.sql, write a query that, in streaming mode, inner joins orders and users on user_id and produces orders (count) and amount (sum of amount) per tier, and save the output to /root/flink/joins/regular.out.
A regular join uses only equality, with no time condition. u41–u44 are users that are not in users, so they drop out of the inner join. Check in the result whether the orders of a user who signed up in the afternoon (signup_time later than the order) are also joined — a regular join finds every partner, past and future, regardless of arrival order.
Interval join — attach only deliveries within 2 hours
In /root/flink/joins/interval.sql, write an interval join that joins orders and shipments on order_id but keeps only those whose ship_time is within 2 hours of the order time (both ends inclusive), and save order_id, ship_id, and delay_s (delay in seconds, TIMESTAMPDIFF(SECOND, ...)) to /root/flink/joins/interval.out.
An interval join needs one equality and a range that ties the times of both sides. BETWEEN a AND b includes both ends. There are deliveries with a delay of exactly 0 seconds and 7200 seconds, and also one with 7201 seconds. An order that was delivered in two separate shipments becomes two lines.
Outer interval join — orders with no partner in time
In /root/flink/joins/unshipped.sql, write a query that uses orders LEFT JOIN shipments to produce the order_id and order_time of orders that had no delivery at all within 2 hours, and save the output to /root/flink/joins/unshipped.out.
Put the time condition in the ON clause, and pick with WHERE only the rows filled with null because there is no partner. Not only orders with no delivery record at all, but also orders that went out after more than 2 hours, must be included. The null row comes out once after the watermark passes (order time + 2 hours).
Temporal join — the exchange rate at the order time
In /root/flink/joins/temporal.sql, make a versioned view rates_v that keeps only the latest row per currency from rates, join orders with FOR SYSTEM_TIME AS OF o.order_time, and produce order_id, currency, rate, and amount_krw (amount × rate). Save the output to /root/flink/joins/temporal.out.
Since rates is append-only, you cannot put a primary key on it. If you make a view that keeps only rows where ROW_NUMBER() OVER (PARTITION BY currency ORDER BY update_time DESC) is 1, it becomes a versioned view with currency as the primary key and update_time as the event time. With an inner join, orders earlier than the first exchange rate drop out.
The same view with a regular join — past orders are recalculated
In /root/flink/joins/latest.sql, join rates_v from step 5 to orders on currency without FOR SYSTEM_TIME AS OF, produce orders (count) and total_krw (sum of amount × rate) per currency, and save the output to /root/flink/joins/latest.out.
A versioned view emits an update every time the exchange rate changes. A regular join receives that update and recomputes even past orders, so you see -U in the output. The final total differs from the total of the temporal join — think about which exchange rate it was computed with.
Read the state in the execution plans of the three joins
Turn the joins of steps 2, 3, and 5 into INSERTs into a blackhole sink and create /root/flink/joins/regular-plan.json, /root/flink/joins/interval-plan.json, and /root/flink/joins/temporal-plan.json with COMPILE PLAN. (The regular join is orders⋈users, the interval join is orders⋈shipments, and the temporal join is orders⋈rates_v.)
The shape is COMPILE PLAN 'file:///path.json' FOR INSERT INTO sink SELECT ... . If the file already exists, it raises an error, so delete it and run again when extracting anew. After creating them, skim the type and state of the nodes with jq — the point is which join nodes have leftState and rightState and which do not.
Report — what was joined and what dropped out in each join
In /root/flink/joins/report.json, write matched_orders (the number of orders joined with users by the regular join), interval_rows (the number of result rows of the interval join), unshipped_orders (the number of rows in step 4), temporal_rows (the number of result rows of the temporal join), orders_without_rate (all orders − temporal_rows), and usd_gap_krw (the USD total_krw in latest.out − the sum of the USD amount_krw in temporal.out, a decimal).
You can copy all of them from the output you saved earlier. For a table with a changelog, the last +I/+U is the final value. If you pick only rows like '| +I |' with grep and cut the columns with awk -F'|', you can count them. Write integer columns as integers and usd_gap_krw as a number.