TT Lab
Get started
Learn Learning paths Courses

Apache Flink — Running Streams on a Real Engine

Reading state and TTL from execution plans

Continue in TT Lab

Goal

Read from the COMPILE PLAN JSON what state each operator holds and what its TTL is, and confirm where a job-wide TTL and per-operator hints are stamped, and what a window aggregation and the state backend change.

Why it matters

An unbounded GROUP BY and a regular join hold state forever by default, and if you reduce it with a TTL the results can be wrong. A TTL means "after this long since the last update, it will be erased sometime", so you cannot judge it from results, but the plan comes out written deterministically. This plan is precisely the first tool for dealing with state problems in operations. The grader of this lab does not ask the cluster anything; it reads the plan JSON, REST responses, and sql-client output you saved, and compares the results with values it computes directly from the source CSV.

Steps

  1. After flink-up, define orders, payments, and users and the blackhole sinks kv_sink (k STRING, v BIGINT) and win_sink (window_start, window_end, v BIGINT) in /root/flink/state/ddl.sql, and extract the plan of an INSERT that only filters with /root/flink/state/calc.sql into /root/flink/state/calc-plan.json.
  2. Extract the plan of an unbounded GROUP BY that produces SUM(amount) per user with /root/flink/state/agg.sql into /root/flink/state/agg-plan.json (without setting a TTL).
  3. Extract the same plan with SET 'table.exec.state.ttl' = '30 min'; added with /root/flink/state/agg-ttl.sql into /root/flink/state/agg-ttl-plan.json.
  4. With the job default at 30 min, extract the plan of the orders o · payments p join with STATE_TTL('o' = '1d', 'p' = '2h') given, using /root/flink/state/join-hint.sql, into /root/flink/state/join-hint-plan.json.
  5. With the job default at 30 min, extract the plan of the orders o ⋈ payments p ⋈ users u join with STATE_TTL('o' = '1d', 'p' = '2h', 'u' = '7d') given, using /root/flink/state/cascade.sql, into /root/flink/state/cascade-plan.json.
  6. With the job default at 30 min, extract the plan of a 10-minute TUMBLE window aggregation into /root/flink/state/window-plan.json, and run the window_start, window_end, orders, and amount of the same window in streaming and save the output to /root/flink/state/window.out (both in one file, /root/flink/state/window.sql).
  7. With SET 'state.backend.type' = 'rocksdb' and the job name flk-state-rocks, save the output of /root/flink/state/rocks.sql, which produces orders and total per user, to /root/flink/state/rocks.out, that job's /jobs/<jid> to /root/flink/state/rocks-job.json, and /jobs/<jid>/checkpoints/config to /root/flink/state/rocks-ckpt.json.
  8. In /root/flink/state/report.json, write default_ttl, join_left_ttl, join_right_ttl, cascade_second_left_ttl, window_state_entries, state_backend, and rocks_groups.

Notes

A query that only filters has no state

After flink-up, define orders, payments, and users and the blackhole sinks kv_sink and win_sink in /root/flink/state/ddl.sql. In /root/flink/state/calc.sql, write a statement that extracts with COMPILE PLAN 'file:///root/flink/state/calc-plan.json' an INSERT that puts the order_id and amount * 2 of orders with amount > 1000 into kv_sink, and run it with sql-client.sh -i ddl.sql -f calc.sql.

COMPILE PLAN takes an INSERT statement, so it needs a sink to receive it. The blackhole connector discards what it receives. Open the JSON you made with jq and look at the type and state of the nodes — an operator that looks at one row and emits it right away has nothing to remember.

An unbounded GROUP BY — forever by default

In /root/flink/state/agg.sql, extract into /root/flink/state/agg-plan.json the plan of an INSERT that puts the per-user SUM(amount) into kv_sink. Do not set a TTL.

Look at what is in the state list of the aggregation node and what the ttl is. 0 ms means never erase. If the result type of SUM differs from the sink column (BIGINT), match it with CAST.

Set a job-wide TTL

At the very front of /root/flink/state/agg-ttl.sql, put SET 'table.exec.state.ttl' = '30 min'; and extract the plan of the same GROUP BY as step 2 into /root/flink/state/agg-ttl-plan.json.

A SET applies from the statements after it. Look at how the groupAggregateState ttl in the plan changes. This value means 'keep at least this long since the last update', so you cannot tell from the results when it was erased.

A different TTL on each side of a join — the STATE_TTL hint

In /root/flink/state/join-hint.sql, keeping SET 'table.exec.state.ttl' = '30 min';, give /*+ STATE_TTL('o' = '1d', 'p' = '2h') */ to an INSERT into kv_sink that regular-joins orders o and payments p on order_id, and extract the plan into /root/flink/state/join-hint-plan.json.

Write the hint right after SELECT. If you attached aliases to the tables, the hint keys must also be aliases. Look at whether the ttl of leftState and rightState at the join node in the plan comes out different from the job default.

A chained join — the slot the hint does not reach

In /root/flink/state/cascade.sql, with the job default at 30 min, extract into /root/flink/state/cascade-plan.json the plan of an INSERT into kv_sink of orders o ⋈ payments p (order_id) ⋈ users u (user_id) with /*+ STATE_TTL('o' = '1d', 'p' = '2h', 'u' = '7d') */ given.

Two join nodes come out. The one with the smaller id is the first (lower) join. Look at which of the four slots (the left and right of the first join, the left and right of the second join) the three hints attach to, and what the one remaining slot receives.

A window aggregation — it clears itself without a TTL

In one file, /root/flink/state/window.sql, put SET 'table.exec.state.ttl' = '30 min'; and (1) extract into /root/flink/state/window-plan.json the plan of an INSERT that puts the 10-minute TUMBLE window SUM(amount) into win_sink, and (2) run the window_start, window_end, orders (COUNT), and amount (SUM) of the same window as a streaming SELECT. Save the output to /root/flink/state/window.out.

The window TVF is FROM TABLE(TUMBLE(TABLE orders, DESCRIPTOR(ts), INTERVAL '10' MINUTE)), grouped with GROUP BY window_start, window_end. Look at whether the window aggregation node has a state entry even though you set a TTL. The fact that the op of the result is all +I is for the same reason — a window emits once when the watermark passes its end and then clears.

Run the same aggregation with the RocksDB backend

In /root/flink/state/rocks.sql, put SET 'state.backend.type' = 'rocksdb'; and SET 'pipeline.name' = 'flk-state-rocks';, and produce orders (COUNT) and total (SUM(amount)) per user in streaming. Save the output to /root/flink/state/rocks.out, that job's /jobs/<jid> to /root/flink/state/rocks-job.json, and /jobs/<jid>/checkpoints/config to /root/flink/state/rocks-ckpt.json.

Find the job id by name in /jobs/overview. The name of the backend class this job used is written in the state_backend field of the checkpoints/config response. A job run with the default backend shows a different name. The results must be the same regardless of the backend.

Report — the state map in numbers

In /root/flink/state/report.json, write default_ttl (the groupAggregateState ttl in agg-plan), join_left_ttl and join_right_ttl (join-hint-plan), cascade_second_left_ttl (the leftState ttl of the second join in cascade-plan), window_state_entries (the number of state entries in window-plan, an integer), state_backend (the value in rocks-ckpt), and rocks_groups (the number of rows in the final result of rocks.out = the number of keys left in state, an integer).

For the ttl values, just copy the strings written in the plan (for example "2 h"). If you sort the nodes whose type starts with stream-exec-join by id with jq, you can tell the two joins apart. For a node with no state, that field is null.