skip to content

In Flink SQL, when do you enrich orders with an event-time temporal join rather than a lookup join, and what does each cost?

level: middleimportance: must knowfreq 58%

answer

  1. same AS OF clause, different time
  2. a versioned table with a key
  3. watermarks on both sides
  4. query the database per row
  5. a cache trades freshness

basics

~20 s

Use an event-time temporal join (FOR SYSTEM_TIME AS OF o.order_time) on a versioned table when each order needs the value valid at its own time, reproducibly. Use a lookup join (AS OF o.proc_time) to query an external table's current row per order.

solid answer

~50 s

Both use `FOR SYSTEM_TIME AS OF`; the time attribute and the right-hand table decide which join Flink builds. For currency rates, an **event-time temporal join** `JOIN currency_rates FOR SYSTEM_TIME AS OF o.order_time AS r ON o.currency = r.currency` needs a versioned table: a changelog with a `PRIMARY KEY` that appears in the join condition, plus an event-time attribute and watermark. Flink keeps rate versions in state, waits for watermarks on both sides, and gives each order the rate valid at `order_time`, so a replay reproduces the same result. For a customer dimension in a database, a **lookup join** `JOIN customers FOR SYSTEM_TIME AS OF o.proc_time AS c` against a `'connector' = 'jdbc'` table queries the database as each order is processed. It keeps no join state, but a replay sees today's data, later changes never revise emitted rows, and throughput leans on `lookup.cache = 'PARTIAL'`, which trades freshness for fewer queries.

code

sql · 6 lines
sql
SELECT o.order_id, o.amount * r.rate AS amount_eur, c.country
FROM orders AS o
JOIN currency_rates FOR SYSTEM_TIME AS OF o.order_time AS r
  ON o.currency = r.currency
JOIN customers FOR SYSTEM_TIME AS OF o.proc_time AS c
  ON o.customer_id = c.id;

go deeper

for a junior

Recall that both joins use FOR SYSTEM_TIME AS OF, event time for a temporal join and processing time for a lookup join.

for a middle

Explain what a versioned table needs, why the temporal join waits for watermarks, and what the lookup cache options trade.

for a senior

Choose the join from the correctness requirement, predict replay behaviour, and tune or disable the lookup cache knowing how stale or missing rows will surface.

for a principal

Decide where reference data should live: capturing a table's changes into Flink for reproducible as-of joins, or querying the system of record and accepting load and drift.

## One syntax, two joins Flink SQL writes both enrichment joins with the SQL:2011 clause `FOR SYSTEM_TIME AS OF`. What follows `AS OF`, and what backs the right-hand table, decide the operator: | | Event-time temporal join | Lookup join | |---|---|---| | `AS OF` argument | the probe side's event-time attribute (`o.order_time`) | the probe side's processing-time attribute (`o.proc_time`) | | Right side | a **versioned table**: a changelog with a `PRIMARY KEY` and an event-time attribute | a table backed by a **lookup source connector**, such as `'connector' = 'jdbc'` | | Where the reference data lives | in Flink state, as versions per key | in the external system, queried per row | | Version matched | the one valid at the order's own time | whatever the external table holds when the row is processed | | Replaying old orders | reproduces the same result | attaches data current at replay time | | Later changes on the right side | do not revise emitted rows | do not revise emitted rows | ## Event-time temporal join: currency rates ```sql SELECT o.order_id, o.amount * r.rate AS amount_eur FROM orders AS o JOIN currency_rates FOR SYSTEM_TIME AS OF o.order_time AS r ON o.currency = r.currency; ``` Requirements and behaviour: - `currency_rates` must be a **versioned table**: a primary key (here `currency`) and an event-time attribute. Typical sources are an `upsert-kafka` topic, a Debezium-format changelog, or an append-only rates stream turned into a changelog by a deduplication query. - The primary key must appear in the equality condition. - The join is **triggered by watermarks from both sides**: an order waits in state until Flink knows which rate version was valid at `order_time`. - Flink stores rate versions and drops those no longer needed as time advances; old orders are not kept in a time window as in an interval join. - Only inner and left outer temporal joins are supported. The pay-off is correctness: an order placed at 10:00 always gets the 10:00 rate, whether the job processes it live or replays it a month later. ## Lookup join: a customer dimension in a database ```sql SELECT o.order_id, c.country, c.segment FROM orders AS o JOIN customers FOR SYSTEM_TIME AS OF o.proc_time AS c ON o.customer_id = c.id; ``` Here `orders` declares a processing-time column (`proc_time AS PROCTIME()`) and `customers` uses the JDBC connector. For each order Flink queries the database for the matching row at the moment it processes the order. The join needs a mandatory equality predicate, keeps no join state, and suits large dimensions that already live in a database. Its costs: 1. **Latency and database load**: one query per order unless cached; the JDBC connector supports only synchronous lookups. 2. **Non-determinism**: a replay attaches whatever the database holds at replay time. 3. **No revision**: when a customer row changes after an order was enriched, that output row stays as it was. ## The JDBC lookup cache (flink-connector-jdbc 4.1) | Option | Default | Meaning | |---|---|---| | `lookup.cache` | `NONE` | `NONE` or `PARTIAL`; caching is off unless set | | `lookup.partial-cache.max-rows` | none | evict the oldest rows beyond this count | | `lookup.partial-cache.expire-after-write` | none | a row's lifetime after it is cached | | `lookup.partial-cache.expire-after-access` | none | a row's lifetime after its last access | | `lookup.partial-cache.cache-missing-key` | `true` | also cache empty lookup results | | `lookup.max-retries` | `3` | retries when a lookup query fails | Each TaskManager holds its own cache. The cache trades **freshness for throughput**: a cached customer can be stale until its entry expires, and with `cache-missing-key` at its default a customer created after a failed lookup keeps missing until that entry expires or is evicted. Separately, the `LOOKUP` query hint can configure a retry strategy, and asynchronous lookup for connectors that support it. ## Choosing - Need the value **as of the event** (rates, prices, contract terms) and reproducible reprocessing: event-time temporal join. - Need the **current** value of a large dimension already in a database, and can accept replay drift: lookup join with a tuned cache. - Need output that is revised when the reference data changes: neither; that is a regular join, with its unbounded state.

  • Why does an event-time temporal join sometimes stop emitting results for recent orders?
    It is triggered by watermarks from both inputs. An order is held until the combined watermark shows which rate version was valid at its `order_time`; if the rates side's watermark stops advancing, orders accumulate in state and nothing is emitted. Check both inputs' watermarks before suspecting the join.
  • What does lookup.partial-cache.cache-missing-key control, and why does its default matter?
    It decides whether an empty lookup result is cached; the default is `true`. A customer created just after an order's lookup missed stays missing on that TaskManager until the entry expires or is evicted, so orders are dropped by an inner join or get `NULL`s from a left join for that period.
  • What must be true of currency_rates before it can be the versioned side of a temporal join?
    It needs a primary key and an event-time attribute, so Flink can order the versions of each key in time, and the join condition must include that key. An append-only rates stream can be turned into a versioned view by a deduplication query that keeps the latest row per currency.

A temporal join is checking the exchange-rate archive for the day of each receipt; a lookup join is phoning the bank for today's rate whenever a receipt reaches your desk, and jotting answers on a notepad (the cache) that you trust for a while.

saying these in an interview costs you the question

  • A lookup join reprocesses last month's orders with last month's customer data.
  • An event-time temporal join updates emitted rows when a new rate arrives.
  • A lookup join keeps the whole dimension table in Flink state.
  • The JDBC lookup cache is enabled by default.
  • An event-time temporal join works without the primary key in its join condition.