skip to content

In Flink SQL, how does an interval join of ad impressions to clicks differ from a regular join in state and output?

level: middleimportance: must knowfreq 55%

answer

  1. a bound on both time attributes
  2. the watermark retires stored rows
  3. append-only in, append-only out
  4. outer rows wait for the bound
  5. plain timestamps fall back

basics

~20 s

An interval join adds a bound between both inputs' time attributes, so Flink drops stored rows once the watermark passes the bound and emits insert-only output. A regular join keeps both inputs indefinitely, and its outer variants retract and re-emit rows.

solid answer

~50 s

A Flink SQL interval join is an equi-join plus a predicate bounding one input's time attribute against the other's, such as `c.click_time BETWEEN i.impression_time AND i.impression_time + INTERVAL '30' MINUTE`. Because time attributes only move forward, Flink removes impressions and clicks from state once the watermark shows no partner can still arrive, and it emits insert-only rows; a `LEFT JOIN` emits an unmatched impression padded with `NULL`s once its bound has passed. The price is that both inputs must be append-only with time attributes; the planner rejects an updating input. A regular join has no bound: it keeps both inputs until a state TTL expires them, accepts updating inputs, and a regular `LEFT JOIN` first emits the impression with `NULL`s, then retracts that row when a click arrives, so it needs a sink that can apply retractions.

code

sql · 6 lines
sql
SELECT i.impression_id, i.ad_id, c.click_id, c.click_time
FROM impressions i
LEFT JOIN clicks c
  ON i.ad_id = c.ad_id AND i.user_id = c.user_id
 AND c.click_time BETWEEN i.impression_time
                      AND i.impression_time + INTERVAL '30' MINUTE;

go deeper

for a junior

Know the shape of an interval join: an equality on the key plus a BETWEEN bounding one time attribute against the other.

for a middle

Explain how the watermark lets the operator drop rows, why its output is insert-only, and how a regular LEFT JOIN retracts its NULL-padded rows instead.

for a senior

Verify with EXPLAIN that a supposed interval join did not fall back to a regular join, and handle updating inputs that make an interval join impossible.

for a principal

Treat the match bound as a business contract: agree with stakeholders how late a click may count, since that bound fixes both state size and what counts as a conversion.

## The two joins side by side | | Regular join | Interval join | |---|---|---| | Condition | equi-join, no time bound | equi-join plus a bound between both inputs' time attributes | | Inputs accepted | append-only or updating | append-only only, with time attributes | | State held | every row of both inputs, until a TTL (default: never) | only rows whose bound has not yet passed | | Inner join output | forwards the inputs' change kinds | insert-only | | Outer join output | may retract and re-emit rows | insert-only; `NULL`-padded rows once the bound passes | ## Writing the interval join ```sql SELECT i.impression_id, i.ad_id, c.click_id, c.click_time FROM impressions i LEFT JOIN clicks c ON i.ad_id = c.ad_id AND i.user_id = c.user_id AND c.click_time BETWEEN i.impression_time AND i.impression_time + INTERVAL '30' MINUTE; ``` Flink plans an **interval join** when the condition has at least one equality predicate and a predicate that bounds time **on both sides**, using time attributes of the same kind (both event time or both processing time): two range comparisons, a `BETWEEN`, or an equality of the two time attributes. A one-sided comparison such as `c.click_time > i.impression_time` leaves the future open and does not qualify. ## How it cleans up Time attributes advance quasi-monotonically, so the operator's watermark tells Flink when a stored row can no longer find a partner: - an impression can be dropped once the watermark is past `impression_time + 30 minutes`, because no click inside its bound can still arrive; - a click can be dropped once the watermark is past `click_time`, because no impression early enough to match it can still arrive; - nothing the bound allows is lost, so there is no TTL trade-off to make. The state held is roughly the rows inside the bound, not the whole history. Two practical consequences follow. First, only outer rows wait: a matched pair is emitted as soon as its partner arrives, while an unmatched impression waits for its whole bound plus the watermark delay before its `NULL`-padded row appears. Second, widening the bound widens both the retained state and that wait, so the bound should come from how late a click may still count, not from a guess. ## What each join emits The interval join consumes and produces **insert-only** rows. With `LEFT JOIN`, an impression that never gets a click is emitted once, padded with `NULL` click columns, after its bound has expired; a matched impression is emitted as a joined row when the click arrives. Neither is revised later, so an insert-only sink can take the result. A regular `LEFT JOIN` cannot wait for anything. It emits the impression with `NULL`s immediately; when a click arrives later it **retracts** that row and inserts the joined one. That changelog needs a sink able to apply retractions or upserts, and the planner refuses to write it to an insert-only sink. ## When you do not get an interval join 1. **Plain timestamps.** If the compared columns are ordinary `TIMESTAMP`s rather than time attributes, the planner treats the predicate as a normal filter and plans a **regular join**, which keeps both inputs indefinitely. 2. **Updating input.** If an input is a changelog, for example a table read with the `upsert-kafka` connector or a `'format' = 'debezium-json'` topic, the interval join cannot consume it and the planner rejects the query. 3. **Mixed time kinds.** A predicate relating an event-time attribute to a processing-time attribute is not a valid interval bound. `EXPLAIN` shows `IntervalJoin` when you got one and `Join` when you did not; check it whenever a time predicate is meant to bound state. ## Choosing between them - Use an **interval join** when the business defines a match window between two event streams: a click counts if it lands within 30 minutes of its impression. - Use a **regular join** when there is no meaningful bound or an input is updating, and pair it with state retention. - Use a **window join** when matches should be confined to the same window TVF window rather than a per-row range. - Use a **temporal or lookup join** when one side is reference data rather than an event stream.

  • What happens if the BETWEEN compares two plain TIMESTAMP columns rather than time attributes?
    Flink does not plan an interval join. The planner treats a predicate as an interval bound only when both sides reference time attributes of the same kind; otherwise the condition becomes an ordinary filter on a regular join, which keeps both inputs indefinitely. `EXPLAIN` shows `Join` instead of `IntervalJoin`.
  • Why can an interval join not take a table read with the upsert-kafka connector as an input?
    An upsert-kafka table is a changelog with updates and deletes, and the interval join only consumes insert-only input, so the planner rejects the query. Use a regular join with state retention, or, if that side is reference data keyed by a primary key, an event-time temporal join.
  • How does an interval join differ from a window join in Flink SQL?
    A window join matches rows that fall into the same window TVF window, with `L.window_start = R.window_start AND L.window_end = R.window_end`, and emits only when the window ends. An interval join uses a range relative to each row, so an impression at 10:59 still matches a click at 11:01 even though they sit in different hourly windows.

saying these in an interview costs you the question

  • A regular join with a timestamp filter is cleaned up automatically.
  • An interval join emits an update when a late click finally matches.
  • An interval join can read an upsert or CDC table as either input.
  • A regular LEFT JOIN emits each impression once, after the match is known.
  • Any BETWEEN on two timestamp columns makes Flink plan an interval join.