TT Lab
Get started
Learn Learning paths Courses

Apache Flink — Running Streams on a Real Engine

Joining orders three ways

Continue in TT Lab

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

  1. Start the cluster with flink-up and 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 names tbl and n).
  2. In /root/flink/joins/regular.sql, write a streaming query that joins orders and users on user_id and produces orders (count) and amount (sum) by tier, and save the output to /root/flink/joins/regular.out.
  3. 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, and delay_s (seconds) to /root/flink/joins/interval.out.
  4. In /root/flink/joins/unshipped.sql, write an outer interval join that produces the order_id and order_time of orders that were not delivered within 2 hours, and save the output to /root/flink/joins/unshipped.out.
  5. Make a versioned view rates_v from rates, and save the output of /root/flink/joins/temporal.sql, a temporal join that attaches the exchange rate at the order time and produces order_id, currency, rate, and amount_krw (amount × rate), to /root/flink/joins/temporal.out.
  6. Save the output of /root/flink/joins/latest.sql, which joins the same rates_v without FOR SYSTEM_TIME AS OF and produces orders and total_krw per currency, to /root/flink/joins/latest.out.
  7. Extract the execution plans of the three joins (steps 2, 3, and 5) with COMPILE PLAN to /root/flink/joins/regular-plan.json, /root/flink/joins/interval-plan.json, and /root/flink/joins/temporal-plan.json.
  8. In /root/flink/joins/report.json, write matched_orders, interval_rows, unshipped_orders, temporal_rows, orders_without_rate, and usd_gap_krw.

Notes

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.