Apache Flink — Running Streams on a Real Engine
Reading state and TTL from execution plans
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
- After
flink-up, define orders, payments, and users and the blackhole sinkskv_sink(k STRING, v BIGINT) andwin_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. - 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). - 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. - 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. - 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. - 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, andamountof the same window in streaming and save the output to /root/flink/state/window.out (both in one file, /root/flink/state/window.sql). - With
SET 'state.backend.type' = 'rocksdb'and the job nameflk-state-rocks, save the output of /root/flink/state/rocks.sql, which producesordersandtotalper user, to /root/flink/state/rocks.out, that job's/jobs/<jid>to /root/flink/state/rocks-job.json, and/jobs/<jid>/checkpoints/configto /root/flink/state/rocks-ckpt.json. - In /root/flink/state/report.json, write
default_ttl,join_left_ttl,join_right_ttl,cascade_second_left_ttl,window_state_entries,state_backend, androcks_groups.
Notes
- Sources (CSVs with no header, times in seconds, ts ascending):
state_orders.csv=order_id, user_id, amount, ts·state_payments.csv=pay_id, order_id, pay_method, ts·state_users.csv=user_id, region. All of them are in/opt/lab/fixtures/data/. Put watermarks on thetsof orders and payments. - Extracting a plan:
COMPILE PLAN 'file:///root/flink/state/이름.json' FOR INSERT INTO 싱크 SELECT ...;(the placeholders stand for the file name and the sink) — it does not run the job. If a file exists at the same path, it raises an error instead of overwriting, so delete it first when extracting again. - Skimming a plan:
jq -c '.nodes[] | {id, type, state}' 파일.json(the placeholder stands for the file) — the order of node ids is from the bottom (source) to the top (sink). - A common mistake: if you attached an alias to a table as the docs do, the hint key must also be the alias.
methodis an SQL reserved word, so using it as a column name makes CREATE fail (pay_method). - Official docs: Hints — STATE_TTL · Configuration — table.exec.state.ttl · State Backends · Group Aggregation · Determinism
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.