TT Lab
开始
学习 学习路径 课程

Apache Flink — 用真正的引擎跑流处理

从执行计划中读取状态与 TTL

在 TT Lab 中继续学习

目标

通过 COMPILE PLAN JSON 读出每个算子持有什么状态、TTL 是多少,并确认整个作业的 TTL 和算子级提示分别落在哪里,以及窗口聚合和状态后端会改变什么。

为什么重要

无尽的 GROUP BY 和常规连接默认会永远持有状态,用 TTL 缩减则可能使结果出错。TTL 的含义是“距上次更新过了这么久之后,会在某个时候清除”,所以无法从结果上判断,而计划是以确定的方式写出来的。在生产环境中处理状态问题,第一个工具就是这份计划。本实验的评分器不会查询集群,而是读取你保存的计划 JSON · REST 响应 · sql-client 输出,并直接用源 CSV 计算结果进行核对。

步骤

  1. 运行 flink-up 后,在 /root/flink/state/ddl.sql 中定义 orders · payments · users 和 blackhole Sink kv_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。
  2. 用 /root/flink/state/agg.sql 把按用户输出 SUM(amount) 的无尽 GROUP BY 的计划导出到 /root/flink/state/agg-plan.json(不设置 TTL)。
  3. 用 /root/flink/state/agg-ttl.sql 把加入了 SET 'table.exec.state.ttl' = '30 min'; 的同一份计划导出到 /root/flink/state/agg-ttl-plan.json。
  4. 保持作业默认值 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。
  5. 在作业默认值 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。
  6. 保持作业默认值 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 一个文件里)。
  7. 把设置了 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。
  8. 在 /root/flink/state/report.json 中写入 default_ttl、join_left_ttl、join_right_ttl、cascade_second_left_ttl、window_state_entries、state_backend、rocks_groups。

参考

只做过滤的查询没有状态

运行 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。