skip to content

A Flink SQL job regular-joins orders to shipments on order_id and its state grows until checkpoints fail. How do you bound it, and what does that cost?

level: seniorimportance: must knowfreq 50%

answer

  1. no time bound in the query
  2. a default of zero
  3. pipeline key versus query hint
  4. per-input retention by alias
  5. forgotten keys lose late matches

basics

~20 s

A regular join keeps both inputs in state, and table.exec.state.ttl defaults to 0, never clean up. Use an interval join if a time bound exists; otherwise set that key or a STATE_TTL hint, accepting that matches after expiry are silently lost.

solid answer

~50 s

A regular join such as `FROM orders o JOIN shipments s ON o.order_id = s.order_id` must be able to match a row arriving at any future time, so Flink keeps every row of both inputs in state, and `table.exec.state.ttl` defaults to `0`, which means state is never cleaned up. If the business gives a real bound (orders ship within seven days), rewrite it as an interval join on the time attributes; Flink then cleans state by watermark with no loss the bound allows. Otherwise set idle-state retention: `SET 'table.exec.state.ttl' = '7 d'` pipeline-wide, or `/*+ STATE_TTL('o' = '7d', 's' = '1d') */` to give each join input its own value. The cost is correctness: a key idle longer than its TTL is forgotten, so a late shipment finds no order and an inner join drops it without any error.

code

sql · 6 lines
sql
SET 'table.exec.state.ttl' = '7 d';

SELECT /*+ STATE_TTL('o' = '7d', 's' = '1d') */
       o.order_id, o.amount, s.carrier
FROM orders o
JOIN shipments s ON o.order_id = s.order_id;

go deeper

for a junior

Remember that a regular Flink SQL join keeps both inputs in state and that table.exec.state.ttl defaults to 0, meaning never clean up.

for a middle

Explain how the pipeline key and the STATE_TTL hint differ in scope, why TTL is a minimum rather than an exact deadline, and which operators clean up by time instead.

for a senior

Walk through diagnosing the growth, prefer an interval join when a bound exists, choose per-input TTLs from measured gaps, and name exactly which results expiry will silently change.

for a principal

Frame TTL as a product decision: the business must sign off on which late matches may be dropped, because the loss is silent and shows up only in downstream numbers.

## Why the regular join's state never shrinks A **regular join** in Flink SQL is an equi-join with no time bound and no `FOR SYSTEM_TIME AS OF` clause: ```sql SELECT o.order_id, o.amount, s.carrier FROM orders o JOIN shipments s ON o.order_id = s.order_id; ``` On unbounded inputs, a shipment may arrive for an order seen months ago. To stay correct for any arrival order, the join operator keeps **every row of both inputs** in keyed state, indexed by `order_id`. Nothing in the query tells Flink when a row can no longer match, and the pipeline-wide idle-state retention, `table.exec.state.ttl`, defaults to `0`, which Flink defines as never cleaning up. State therefore grows with every order and shipment ever seen, checkpoints grow with it, and eventually they time out or the job runs short of disk or memory. ## Option 1: give the join a time bound If the business can state a bound — every order ships within seven days — rewrite the query as an **interval join** on the two time attributes: ```sql SELECT o.order_id, o.amount, s.carrier FROM orders o JOIN shipments s ON o.order_id = s.order_id AND s.ship_time BETWEEN o.order_time AND o.order_time + INTERVAL '7' DAY; ``` Flink then removes rows as the watermark passes the bound, and nothing the bound allows is lost. This needs append-only inputs with time attributes, so it is not always available. ## Option 2: `table.exec.state.ttl` `SET 'table.exec.state.ttl' = '7 d';` sets **idle state retention** for the stateful operators of the pipeline, including regular joins, continuous `GROUP BY`, Top-N and deduplication. Its semantics, as the option itself defines them: - It is a **minimum**: state that has not been updated for less than the TTL is kept; state idle for longer is removed at some point afterwards, not at an exact instant. - The default `0` means state is never cleaned up. - The bookkeeping for expiry adds some overhead of its own. ## Option 3: the `STATE_TTL` hint A pipeline-wide value is blunt: orders may need a week while shipments need a day, and a `GROUP BY` elsewhere in the job may need no expiry at all. The `STATE_TTL` query hint sets **operator-level** retention for regular joins and group aggregations, overriding the pipeline value: ```sql SELECT /*+ STATE_TTL('o' = '7d', 's' = '1d') */ o.order_id, s.carrier FROM orders o JOIN shipments s ON o.order_id = s.order_id; ``` Keys are table names or aliases; once a table has an alias, the hint must use the alias. The hint applies only to its own query block. In a chain of joins, the values cover both inputs of the first join and the right input of each later one, while the left input of later joins falls back to `table.exec.state.ttl` unless you split the query into views. ## What TTL costs TTL buys bounded state by forgetting keys. What forgetting means depends on the operator: | Operator | Effect of an expired key | |---|---| | Inner regular join | A late partner finds nothing and **no row is produced**; no error is raised. | | Left outer regular join | An order already emitted with `NULL` shipment columns stays that way when its shipment arrives after the order's state expired. | | Continuous `GROUP BY` | The key's aggregate **restarts from zero**, as if never seen. | | Deduplication | A duplicate arriving after expiry is treated as a first occurrence. | So the TTL must exceed the longest gap you expect between matching rows, and the business must accept that anything beyond it is dropped quietly. ## Diagnosis checklist 1. Confirm the growth is in the join: per-operator checkpoint sizes in the web UI, and `EXPLAIN` showing a `Join` rather than an `IntervalJoin`. 2. Check `table.exec.state.ttl`; `0` means unbounded. 3. Ask whether a real time bound exists; if so, prefer an interval join. 4. Ask whether it is really enrichment against reference data; a temporal or lookup join holds far less state. 5. Otherwise set per-input `STATE_TTL` values from measured gaps, and watch the business metric that expiry would distort. 6. For several chained regular joins sharing a key, consider the multi-join operator introduced in Flink 2.1, which avoids storing intermediate join results.

  • Which Flink SQL operators need idle-state retention, and which clean themselves up?
    Operators with no time bound need it: regular joins, continuous `GROUP BY`, Top-N and deduplication keep per-key state until `table.exec.state.ttl` or a `STATE_TTL` hint expires it. Window TVF aggregations and interval joins are bounded by time: they drop state once the watermark passes the window end or the join bound, so TTL is not what keeps them small.
  • Does table.exec.state.ttl remove state exactly when the configured time elapses?
    No. The value is a minimum: state idle for less than the TTL is always kept, and state idle for longer is removed at some later point. Size falls lazily, so capacity planning must assume expired entries linger for a while, and expiry bookkeeping itself costs some overhead.
  • Why would you use the STATE_TTL hint rather than only the pipeline key?
    The pipeline key applies one value to every stateful operator, so a week for orders also forces a week on shipments and on any unrelated aggregation. The hint sets retention per input of a regular join or group aggregation, keyed by table name or alias, so each side keeps only as long as its own matches need.

saying these in an interview costs you the question

  • table.exec.state.ttl defaults to one day, so a regular join cleans itself up.
  • Setting a state TTL on a regular join costs nothing in correctness.
  • An inner join whose partner expired raises an error you can alert on.
  • The STATE_TTL hint also sets retention for window TVF aggregations.
  • Moving to a disk-based state backend fixes a join whose state grows forever.