Apache Flink — Running Streams on a Real Engine
Stream joins — what is remembered, and for how long
In one line
Stream joins are divided by which side's rows are remembered and for how long. A regular join with no time condition holds both sides forever, an interval join discards rows once the time range has passed, and an event-time temporal join holds only the versions of the right side and joins a left row with exactly one "version as of that time".
Why this was needed
In batch, a join is easy. Both tables are already complete, so you build a hash table from one side, scan the other side, and you are done. In a stream, the two tables do not end. At the moment an order comes in, the delivery of that order does not exist yet, and the user information may have arrived long ago or may arrive long after. So the engine has to pile up "the rows on the other side that could be a partner of the row that just arrived" somewhere, and that pile is the state.
The problem is when to discard. If you hold forever on to a row whose partner may never arrive, the state grows endlessly. But if you discard at random times, you miss a partner that arrives late. This is why Flink SQL divides joins into several kinds — depending on what the query promises about time, what the engine can safely discard differs.
How it works
A regular join is the most permissive. In the wording of the official docs, when a new row arrives on one side, it is matched against all rows of the other side, past and future. So even the morning order of a user who signed up in the afternoon gets joined — because it does not look at time. The price is also written in the docs. Both inputs have to be kept in state forever. You can reduce it with state TTL, but then the results can be wrong. If you extract the execution plan as JSON (COMPILE PLAN), the join node shows two states, leftState and rightState, with a TTL of 0 ms (never erased).
An interval join requires one equality condition and a range that ties the times of both sides. s.ship_time BETWEEN o.order_time AND o.order_time + INTERVAL '2' HOUR is an example. The inputs have to be append-only tables with a time attribute. Since time attributes increase almost monotonically, once the watermark passes (order time + 2 hours), it is settled that no more deliveries that could pair with that order will come, and it is erased from state. There is also a boundary confirmed by measurement. BETWEEN includes both ends, so a delivery with a delay of 0 seconds and one of exactly 7200 seconds are joined, while one of 7201 seconds is not. If you change it to a LEFT JOIN and put the time condition in ON, an order that found no partner by the time the window closed comes out once with nulls. The result is append-only to the end.
An event-time temporal join joins one left (order) row with exactly one "version that was valid at that time" of the right versioned table. The syntax is FOR SYSTEM_TIME AS OF o.order_time from SQL:2011. To be a versioned table, it needs a primary key and an event-time attribute. You cannot put a primary key on an append-only source such as an exchange rate file, and the docs give a trick for this. If you create a deduplication view with ROW_NUMBER() OVER (PARTITION BY currency ORDER BY update_time DESC) = 1 per currency, the optimizer infers currency as the primary key and uses it as a versioned view.
CREATE TEMPORARY VIEW rates_v AS
SELECT currency, rate, update_time FROM (
SELECT *, ROW_NUMBER() OVER (PARTITION BY currency ORDER BY update_time DESC) AS rn FROM rates)
WHERE rn = 1;
SELECT o.order_id, r.rate
FROM orders o JOIN rates_v FOR SYSTEM_TIME AS OF o.order_time AS r
ON o.currency = r.currency;
In the measured results, an order was joined with the latest version among those with update_time <= order_time (if the exchange rate changes in the same second as the order, the new rate). An order earlier than the first exchange rate had no version to join and dropped out of the inner join. As the docs say, this join is triggered by the watermarks of both sides, and even if the right side changes later, it does not fix the results it already emitted. Old versions are erased from state once they are no longer needed. If you look at the plan JSON, the temporal join node and the interval join node have no state entry with a TTL at all — these two clean up state by time, not by TTL.
If you join the same versioned view without FOR SYSTEM_TIME AS OF, it is just a regular join. Every time the exchange rate changes, past orders are recomputed with the new rate and retractions (-U) and updates (+U) pour out, and the final result becomes the value converted at the last exchange rate.
| Join | What it remembers | Does it fix results | Basis for clearing state |
|---|---|---|---|
| Regular | Everything on both sides | Yes (if the input is updating) | Only TTL (may lose correctness) |
| Interval | Rows within the time range | No | The watermark passes the end of the range |
| Event-time temporal | The needed versions of the right side | No | Versions no longer needed after the watermark passes |
What it looks like in the field
The most common accident is "I only attached product information to orders, but the state has been growing for months". The cause is almost always a regular join. The product table is small, but the order side piles up forever. If the dimension information only needs to be the value at that point in time, a temporal join is the right tool.
The second is the sales recalculation accident. If you attach an exchange rate or a price list with a regular join, the amounts of past orders silently change at the moment the price changes. When a report arrives that yesterday's sales on the dashboard differ today, check the join type first. In this lab, if you convert the same USD orders in both ways, the totals really do differ.
The third is the boundary of an interval join. Whether you write "delivered within 2 hours" with < or with BETWEEN decides the delivery that went out at exactly 2 hours. You need to match the definition of the operational metric with the inequality sign in the SQL at least once.
What you will do in the next lab
You define four sources and collect orders by tier with a regular join. You join deliveries within 2 hours with an interval join, and find orders that were not delivered on time with an outer interval join. You make the exchange rates into a versioned view, join the exchange rate at the order time with a temporal join, and see how the total changes when the same view is joined with a regular join. You extract the execution plans of the three joins as JSON to compare their state entries, and write up the numbers as a report.