状态存在哪里、有多少——算子、TTL 与状态后端
一句话总结
状态不是属于整个作业,而是每个算子各有一份。哪个算子持有什么、TTL 是多少,COMPILE PLAN 会用 JSON 准确地写出来。无尽的 GROUP BY 和常规连接默认会永远持有状态,窗口聚合则会在窗口关闭时自行清空。状态后端只决定这些状态放在哪里、以什么形态存放,并不会改变结果。
为什么需要它
流处理作业运行几周之后,几乎一定会遇到的问题是“内存(或磁盘)为什么一直在增长”。官方文档对 GROUP BY 给出了这样的警告:计算流式查询结果所需的状态可能无限增长,其大小取决于分组数量和聚合函数的种类(MIN、MAX 很重,COUNT 很轻)。上一个模块说过,常规连接会永远持有两侧的输入。
所以要设置 TTL(空闲状态保留时间)。但 TTL 并不是免费的。文档的确定性一章(Determinism)指出,用 TTL 清理状态往往是必要的妥协,但可能让结果变得非确定。被清除的键如果稍后又来了行,就会像第一次见到的键一样重新计数。所以必须以算子为单位,清楚地知道并确定“给哪个算子设置多少”。本模块讲的,就是从执行计划中读出这张地图的方法。
工作原理
COMPILE PLAN 'file:///경로.json' FOR INSERT INTO ...(占位符为路径)不会运行作业,而是把优化后的计划写成 JSON。每个节点都有 type,带状态的节点会附有 state 列表。在本实验环境中导出的结果如下。
| 节点(type) | state 项 | 默认 TTL |
|---|---|---|
| stream-exec-calc(过滤、转换) | 无 | — |
| stream-exec-group-aggregate(无尽的 GROUP BY) | groupAggregateState | 0 ms |
| stream-exec-join(常规连接) | leftState · rightState | 0 ms |
| stream-exec-deduplicate · stream-exec-rank | deduplicateState · rankState | 0 ms |
| 区间连接 · 时态连接 · 窗口聚合 | 无 | — |
0 ms 表示“不清除”。文档对配置(table.exec.state.ttl)的说明就是这样:对没有被更新的状态至少保留这么久,超过这段时间没有动静,之后会在某个时候清除。默认值 0 就是让它永远保留。设置 SET 'table.exec.state.ttl' = '30 min' 后,上表中的 0 ms 会全部变成 30 min。表中最后一行的算子不是靠 TTL,而是靠时间(水位线越过窗口末端、区间末端)来清理状态,所以计划中没有 TTL 项。文档也写道:“生命周期短的窗口 GROUP BY 不是问题。”
给整个作业设置同一个值太粗糙了。有的连接,订单要记一天,而支付只需要两小时。所以有了 STATE_TTL 提示(hint)。它只用于常规连接和 GROUP BY,以表名或别名作为键(如果加了别名,就必须用别名)。
SET 'table.exec.state.ttl' = '30 min';
SELECT /*+ STATE_TTL('o' = '1d', 'p' = '2h') */ ...
FROM orders o JOIN payments p ON o.order_id = p.order_id;
-- 계획: leftState "1 d" · rightState "2 h" (잡 기본값보다 힌트가 우선)
当连接相连时还有一条规则。文档指出,三个提示会被解释到第一个连接的左侧、右侧和第二个连接的右侧,第二个连接的左侧(第一个连接的结果)来自作业配置。实际上,给 STATE_TTL('o'='1d','p'='2h','u'='7d') 配上作业默认值 30 min 后导出的计划里,第二个连接是 leftState "30 min" · rightState "7 d"。如果想把那个位置也定下来,就得把第一个连接拆成视图,再单独给出提示。
状态后端决定这些状态放在哪里。根据文档,什么都不指定时用 HashMapStateBackend,把状态作为 Java 堆中的对象存放。EmbeddedRocksDBStateBackend 则把状态以序列化后的字节存放在 TaskManager 本地磁盘的 RocksDB 中。它能装下磁盘那么多的数据,代价是每次读写都要付出序列化成本。文档指出,提供增量检查点的是这个后端,而把状态放在远程存储中的 ForSt 后端仍处于实验阶段。可以为每个作业用 SET 'state.backend.type' = 'rocksdb' 切换,作业用了哪个后端,会记录在 /jobs/<jid>/checkpoints/config 响应的 state_backend 一栏中。同一个 GROUP BY 即使用 RocksDB 运行,结果也不会有一个字符的差别。
在现场相遇的样子
收到“状态在增长”的告警,先导出计划。哪个算子以 TTL 0 ms 持有什么,一目了然。元凶往往是随手写下的无尽 GROUP BY,或是与维度表的常规连接。能否改成窗口聚合、能否改成时态连接,是比 TTL 更先该考虑的解法。
决定设置 TTL 时,先考虑提示,再考虑整个作业的值。整个作业设置 30 分钟,会把保留一天的订单状态也在 30 分钟后清除,让结果出错。反过来,如果只相信提示,而忘了相连的连接中第二个左侧的位置,那里会使用作业默认值(默认 0 = 永远)。在计划 JSON 里用数字确认,是正确的习惯。
更换后端的原因通常是堆不够。这个 Pod 的 TaskManager 任务堆只有 200MiB 出头,所以大的状态放不进 hashmap。改用 RocksDB 可以让它溢出到磁盘,但会变慢——这是用大小换速度的决定。
下一项实验要做什么
定义三个 Source 和一个丢弃数据的 Sink,在只做过滤的查询的计划中确认没有状态。读出无尽 GROUP BY 的默认 TTL,再设置整个作业的 TTL,观察值的变化。用提示给连接的两侧设置不同的 TTL,并在相连的连接中找出提示够不到的位置。确认窗口聚合的计划和结果,再用 RocksDB 后端运行同一个聚合,最后把所有数字整理成报告。