从执行计划中读取状态与 TTL
目标
通过 COMPILE PLAN JSON 读出每个算子持有什么状态、TTL 是多少,并确认整个作业的 TTL 和算子级提示分别落在哪里,以及窗口聚合和状态后端会改变什么。
为什么重要
无尽的 GROUP BY 和常规连接默认会永远持有状态,用 TTL 缩减则可能使结果出错。TTL 的含义是“距上次更新过了这么久之后,会在某个时候清除”,所以无法从结果上判断,而计划是以确定的方式写出来的。在生产环境中处理状态问题,第一个工具就是这份计划。本实验的评分器不会查询集群,而是读取你保存的计划 JSON · REST 响应 · sql-client 输出,并直接用源 CSV 计算结果进行核对。
步骤
- 运行
flink-up后,在 /root/flink/state/ddl.sql 中定义 orders · payments · users 和 blackhole Sinkkv_sink(k STRING, v BIGINT)、win_sink(window_start, window_end, v BIGINT),并用 /root/flink/state/calc.sql 把只做过滤的 INSERT 的计划导出到 /root/flink/state/calc-plan.json。 - 用 /root/flink/state/agg.sql 把按用户输出
SUM(amount)的无尽 GROUP BY 的计划导出到 /root/flink/state/agg-plan.json(不设置 TTL)。 - 用 /root/flink/state/agg-ttl.sql 把加入了
SET 'table.exec.state.ttl' = '30 min';的同一份计划导出到 /root/flink/state/agg-ttl-plan.json。 - 保持作业默认值 30 min,用 /root/flink/state/join-hint.sql 把给 orders o · payments p 的连接加上
STATE_TTL('o' = '1d', 'p' = '2h')的计划导出到 /root/flink/state/join-hint-plan.json。 - 在作业默认值 30 min 下,用 /root/flink/state/cascade.sql 把 orders o ⋈ payments p ⋈ users u 的连接加上
STATE_TTL('o' = '1d', 'p' = '2h', 'u' = '7d')的计划导出到 /root/flink/state/cascade-plan.json。 - 保持作业默认值 30 min,把 10 分钟 TUMBLE 窗口聚合的计划导出到 /root/flink/state/window-plan.json,并以流式方式运行同一窗口的
window_start、window_end、orders、amount,把输出保存到 /root/flink/state/window.out(两者都在 /root/flink/state/window.sql 一个文件里)。 - 把设置了
SET 'state.backend.type' = 'rocksdb'和作业名称flk-state-rocks、按用户输出orders、total的 /root/flink/state/rocks.sql 的输出保存到 /root/flink/state/rocks.out,把该作业的/jobs/<jid>保存到 /root/flink/state/rocks-job.json,把/jobs/<jid>/checkpoints/config保存到 /root/flink/state/rocks-ckpt.json。 - 在 /root/flink/state/report.json 中写入
default_ttl、join_left_ttl、join_right_ttl、cascade_second_left_ttl、window_state_entries、state_backend、rocks_groups。
参考
- 源数据(没有表头的 CSV,时间以秒为单位,ts 升序):
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。都在/opt/lab/fixtures/data/中。请在 orders、payments 的ts上设置水位线。 - 导出计划:
COMPILE PLAN 'file:///root/flink/state/이름.json' FOR INSERT INTO 싱크 SELECT ...;(占位符依次为计划文件名与 Sink 名称)——不会运行作业。同一路径下已有文件时不会覆盖而是报错,所以重新导出前要先删除。 - 浏览计划:
jq -c '.nodes[] | {id, type, state}' 파일.json(占位符为文件名)——节点 id 的顺序就是从下方(Source)到上方(Sink)的顺序。 - 常见错误:如果按文档给表加了别名,提示的键也必须是别名。
method是 SQL 保留字,用作列名会导致 CREATE 失败(请用pay_method)。 - 官方文档:Hints — STATE_TTL · Configuration — table.exec.state.ttl · State Backends · Group Aggregation · Determinism
只做过滤的查询没有状态
运行 flink-up 后,在 /root/flink/state/ddl.sql 中定义 orders · payments · users 和 blackhole Sink kv_sink、win_sink。在 /root/flink/state/calc.sql 中写出用 COMPILE PLAN 'file:///root/flink/state/calc-plan.json' 导出 INSERT 计划的语句,该 INSERT 把 amount > 1000 的订单的 order_id 和 amount * 2 写入 kv_sink,并用 sql-client.sh -i ddl.sql -f calc.sql 运行。
COMPILE PLAN 接受 INSERT 语句,所以需要一个接收的 Sink。blackhole 连接器会把收到的内容丢弃。用 jq 打开生成的 JSON,看 nodes 的 type 和 state——看一行就立刻输出的算子,没有需要记住的东西。
无尽的 GROUP BY——默认永久保留
用 /root/flink/state/agg.sql 把按用户的 SUM(amount) 写入 kv_sink 的 INSERT 的计划导出到 /root/flink/state/agg-plan.json。不设置 TTL。
看聚合节点的 state 列表里有什么、ttl 是多少。0 ms 表示不清除。如果 SUM 的结果类型与 Sink 列(BIGINT)不同,用 CAST 对齐。
设置整个作业的 TTL
在 /root/flink/state/agg-ttl.sql 的最前面写入 SET 'table.exec.state.ttl' = '30 min';,把与第 2 步相同的 GROUP BY 的计划导出到 /root/flink/state/agg-ttl-plan.json。
SET 从其后的语句开始生效。看计划中 groupAggregateState 的 ttl 变成了什么。这个值的意思是“距上次更新至少保留这么久”,所以无法从结果上得知它是什么时候被清除的。
连接两侧设置不同的 TTL——STATE_TTL 提示
在 /root/flink/state/join-hint.sql 中保持 SET 'table.exec.state.ttl' = '30 min';,对把 orders o 与 payments p 按 order_id 做常规连接并写入 kv_sink 的 INSERT 加上 /*+ STATE_TTL('o' = '1d', 'p' = '2h') */,并把计划导出到 /root/flink/state/join-hint-plan.json。
提示写在 SELECT 之后。如果给表加了别名,提示的键也必须是别名。看计划的连接节点中 leftState · rightState 的 ttl 是否与作业默认值不同。
相连的连接——提示够不到的位置
在 /root/flink/state/cascade.sql 中设置作业默认值 30 min,把对 orders o ⋈ payments p(order_id)⋈ users u(user_id)写入 kv_sink 的 INSERT 加上 /*+ STATE_TTL('o' = '1d', 'p' = '2h', 'u' = '7d') */ 的计划导出到 /root/flink/state/cascade-plan.json。
会出现两个连接节点。id 较小的一个是先(在下方)连接的。看三个提示落在四个位置(第一个连接的左侧、右侧,第二个连接的左侧、右侧)中的哪里,剩下的一个位置接收的是什么。
窗口聚合——不靠 TTL,自行清空
在 /root/flink/state/window.sql 一个文件里写入 SET 'table.exec.state.ttl' = '30 min';,(1) 把 10 分钟 TUMBLE 窗口的 SUM(amount) 写入 win_sink 的计划导出到 /root/flink/state/window-plan.json,(2) 再用流式 SELECT 运行同一窗口的 window_start、window_end、orders(COUNT)、amount(SUM)。输出保存到 /root/flink/state/window.out。
窗口 TVF 是 FROM TABLE(TUMBLE(TABLE orders, DESCRIPTOR(ts), INTERVAL '10' MINUTE)),并用 GROUP BY window_start, window_end 分组。看设置了 TTL 之后,窗口聚合节点里是否有 state 项。结果的 op 全是 +I 也是同样的原因——窗口会在水位线越过末端时输出一次并清空。
用 RocksDB 后端运行同一个聚合
在 /root/flink/state/rocks.sql 中写入 SET 'state.backend.type' = 'rocksdb';、SET 'pipeline.name' = 'flk-state-rocks';,并以流式方式输出按用户的 orders(COUNT)、total(SUM(amount))。输出保存到 /root/flink/state/rocks.out,该作业的 /jobs/<jid> 保存到 /root/flink/state/rocks-job.json,/jobs/<jid>/checkpoints/config 保存到 /root/flink/state/rocks-ckpt.json。
作业 ID 在 /jobs/overview 中按名称查找。checkpoints/config 响应的 state_backend 一栏会写明这个作业使用的后端类名。用默认后端运行的作业,会出现另一个名称。结果必须与后端无关,保持相同。
报告——把状态地图写成数字
在 /root/flink/state/report.json 中写入 default_ttl(agg-plan 中 groupAggregateState 的 ttl)、join_left_ttl 和 join_right_ttl(join-hint-plan)、cascade_second_left_ttl(cascade-plan 中第二个连接的 leftState 的 ttl)、window_state_entries(window-plan 中 state 项的个数,整数)、state_backend(rocks-ckpt 中的值)、rocks_groups(rocks.out 最终结果的行数 = 状态中剩余的键数,整数)。
ttl 值直接抄计划中写的字符串(例如 "2 h")即可。用 jq 把 type 以 stream-exec-join 开头的节点按 id 排序,就能区分两个连接。没有 state 的节点,该栏是 null。