ARTICLE · INTELLIGENCE

战地情报 · 详情页

来自尧图项目组的一线实战观察与深度解析

Flink 流式概念:Table API SQL 流批统一下的状态管理、状态 TTL 与算子级生命周期配置

Flink 流式概念:Table API  SQL 流批统一下的状态管理、状态 TTL 与算子级生命周期配置 大数据流处理批处理数据工程【免费下载链接】flink项目地址https://gitcode.com/gh_mirrors/fli/flink点击查看免费下载Flink 的 Table API 与 SQL 是一套流批统一的声明式 API在有限的批式输入和无限的流式输入下具备相同的语义。由于关系代数与 SQL 最初是为批处理设计的流式场景下的状态管理成为理解与调优的关键。本文以docs/content.zh/docs/dev/table/concepts/overview.md为核心系统讲解流式表程序的状态使用方式、空闲状态维持时间State TTL的三种配置途径、从 Flink v1.18 起支持的算子级状态 TTL含 CompiledPlan 完整实操并延伸状态化更新与演化等进阶话题帮助你掌握状态维度下的流式 SQL 生产实践。流批统一状态是流式表程序的灵魂Flink 的 Table API 与 SQL 在批式与流式输入下共享同一套语义。区别在于批式查询天然拥有有限的输入集可以一次性完成计算而流式查询面对的是无限的数据流必须以**连续查询Continuous Query**的形式持续运行因此必须依赖状态来保存跨时间维度的中间结果。一个流模式下运行的表程序Table program可以完整利用 Flink 作为有状态流处理器的能力配置不同的 state backend如 RocksDB、Heap以适配不同规模的状态存储需求配置多种 checkpoint 选项以满足不同的容错与恢复需求对正在运行的 Table API SQL 管道生成 savepoint并在之后用其恢复应用状态。状态使用声明式管道中的隐式状态由于 Table API SQL 程序是声明式的状态会在哪里、如何被使用并不直接可见。**Planner优化器**负责判断是否需要状态来得到正确的计算结果并尽可能把管道优化成使用更少状态的形式。从概念上讲源表从来不会在状态中被完全保存——实现者处理的是逻辑表即动态表Dynamic Table各算子的状态完全取决于具体用到的操作。显式状态算子Join、聚合与去重包含连接Join、聚合Aggregation或去重Deduplication等操作的语句需要在 Flink 抽象的容错存储内保存中间结果这类算子被称为状态算子。例如对两个表执行普通 Join基于正确的 SQL 语义运行时假设两表会在任意时间点进行匹配因此算子需要保存两个表的全部输入。为了控制状态规模Flink 提供了优化窗口 Join 和时段 Join利用 watermarks 概念即时间属性让过期的数据不再参与匹配从而显著缩小状态。另一个经典例子是词频统计CREATE TABLE doc ( word STRING ) WITH ( connector ... ); CREATE TABLE word_cnt ( word STRING PRIMARY KEY NOT ENFORCED, cnt BIGINT ) WITH ( connector ... ); INSERT INTO word_cnt SELECT word, COUNT(1) AS cnt FROM doc GROUP BY word;这里word是分组的键连续查询为每个观察到的word维护一个中间状态来保存当前词频。由于输入word的值随时间变化且查询持续运行Flink 会为每个word维护一个中间状态总状态量会随着新word的出现不断增长——这正是流式聚合状态下最需要警惕的内存与存储风险点。隐式状态算子SELECT 也可能引入状态形如SELECT ... FROM ... WHERE这种只包含字段映射或过滤器的查询通常是无状态的。但在某些情况下根据输入数据的特征或配置状态算子会被隐式地推导出来输入表是不带UPDATE_BEFORE的更新流详见表到流的转换或配置了table-exec-source-cdc-events-duplicate。下面的例子展示了对 upsert-kafka 源表执行最简单的SELECT *CREATE TABLE upsert_kakfa ( id INT PRIMARY KEY NOT ENFORCED, message STRING ) WITH ( connector upsert-kafka, ... ); SELECT * FROM upsert_kakfa;upsert-kafka 源表的消息类型只包含INSERT、UPDATE_AFTER和DELETE而下游可能要求完整的 changelog包含UPDATE_BEFORE。因此虽然查询本身不包含任何状态计算优化器依然会隐式地推导出一个 ChangelogNormalize 状态算子来生成完整的 changelog。空闲状态维持时间table.exec.state.ttl空闲状态维持时间参数table.exec.state.ttl定义了状态的键在被更新后要保持多长时间才被移除。在上述词频例子中某个word的计数会在配置的时间内未更新时被立刻移除。该配置项的源码定义位于flink-table/flink-table-api-java/src/main/java/org/apache/flink/table/api/config/ExecutionConfigOptions.java键名为table.exec.state.ttl类型为 Duration默认值为0 ms含义是永不清理状态清理状态会引入额外的簿记开销bookkeeping因此默认关闭。移除状态的键之后连续查询会完全忘记它曾经见过这个键如果一条记录带有一个曾被移除状态的键该记录会被当作对应键的第一条记录处理。在词频例子中这意味着cnt会再次从0开始计数——这是配置 TTL 时最容易忽略的语义影响。指定状态生命周期的三种方式从 Flink v1.18 开始Table API SQL 支持多种粒度的状态 TTL 配置方式下表源自原文档总结了它们的适用面与优先级配置方式TableAPI/SQL 支持生效范围优先级SET table.exec.state.ttl ...TableAPI、SQL作业粒度默认情况下所有状态算子都会使用该值控制状态生命周期默认配置可被覆盖SELECT /* STATE_TTL(...) */ ...SQL有限算子粒度当前支持连接和分组聚合算子该值优先作用于相应算子的状态生命周期详见状态生命周期提示修改序列化为 JSON 的 CompiledPlanTableAPI、SQL通用算子粒度可修改任一状态算子的生命周期table.exec.state.ttl与STATE_TTL的值会序列化到 CompiledPlan若作业使用 CompiledPlan 提交最终生效的生命周期由最后一次修改的状态元数据决定STATE_TTL 查询提示STATE_TTL提示以 SQL 注释形式作用于具体算子。从源码flink-table/flink-table-planner/src/main/java/org/apache/flink/table/planner/hint/StateTtlHint.java可以看出其实现要点对于双输入算子如 Join支持形如STATE_TTL(T1 1d, T2 2d)的键值对写法分别指定左右输入的 TTL源码中LEFT_INPUT映射为输入侧 0其余映射为输入侧 1对于单输入算子如分组聚合支持STATE_TTL(T1 2d)形式TTL 值支持d天、h小时等 Flink 时间单位通过TimeUtils.parseDuration解析为毫秒。对应的测试flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/hints/stream/StateTtlHintTest.java覆盖了大量边界场景例如-- 为 Join 的左右输入分别指定 TTL select /* STATE_TTL(T2 2d, T1 1d) */* from T1 join T2 on T1.a1 T2.a2 -- 为分组聚合指定 TTL select /* STATE_TTL(T1 2d) */ count(*) from T1 group by a1测试同时验证了非法用法会被拒绝例如提示选项与输入表名不匹配报错The options of following hints cannot match the name of input tables or views、STATE_TTL()不带任何键值选项报错Invalid STATE_TTL hint, expecting at least one key-value options specified.等情况。配置算子粒度的状态 TTL高级特性注意这是一个需要小心使用的高级特性。它仅适用于作业中使用了多个状态、且每个状态需要不同 TTL 的场景。无状态作业无需关注若作业仅使用一个状态仅需设置作业级 TTL 参数table.exec.state.ttl即可。从 Flink v1.18 开始Table API SQL 支持以每个状态算子的入边数为粒度配置细粒度状态 TTLOneInputStreamOperator单输入可配置一个状态的 TTLTwoInputStreamOperator如双流 Join可分别为左状态和右状态配置 TTL更一般地具有 K 个输入的MultipleInputStreamOperator可以配置 K 个状态 TTL。典型使用场景为双流 Join的左右流配置不同 TTL双流 Join 会生成拥有两条输入边的TwoInputStreamOperator状态算子分别用两个状态保存来自左流和右流的更新在同一作业中为不同的状态计算设置不同 TTL例如一个 ETL 作业先用ROW_NUMBER进行去重再用GROUP BY进行聚合会生成两个拥有单条输入边的OneInputStreamOperator状态算子可为它们分别设置不同的 TTL。需要说明的是基于窗口的操作如窗口连接、窗口聚合、窗口 Top-N 等和 Interval Join 不依赖table.exec.state.ttl控制状态保留因此它们的状态无法在算子级别配置。第一步生成 Compiled Plan配置过程首先使用COMPILE PLAN语句生成一个 JSON 文件它表示序列化后的执行计划。注意COMPILE PLAN不支持查询语句SELECT ... FROM ...只支持INSERT类语句或语句集合。Java 方式TableEnvironment tableEnv TableEnvironment.create(EnvironmentSettings.inStreamingMode()); tableEnv.executeSql( CREATE TABLE orders (order_id BIGINT, order_line_id BIGINT, buyer_id BIGINT, ...)); tableEnv.executeSql( CREATE TABLE line_orders (order_line_id BIGINT, order_status TINYINT, ...)); tableEnv.executeSql( CREATE TABLE enriched_orders (order_id BIGINT, order_line_id BIGINT, order_status TINYINT, ...)); // CompilePlan#writeToFile only supports a local file path, if you need to write to remote filesystem, // please use tableEnv.executeSql(COMPILE PLAN hdfs://path/to/plan.json FOR ...) CompiledPlan compiledPlan tableEnv.compilePlanSql( INSERT INTO enriched_orders \n SELECT a.order_id, a.order_line_id, b.order_status, ... \n FROM orders a JOIN line_orders b ON a.order_line_id b.order_line_id); compiledPlan.writeToFile(/path/to/plan.json);Scala 方式val tableEnv TableEnvironment.create(EnvironmentSettings.inStreamingMode()) tableEnv.executeSql( CREATE TABLE orders (order_id BIGINT, order_line_id BIGINT, buyer_id BIGINT, ...)) tableEnv.executeSql( CREATE TABLE line_orders (order_line_id BIGINT, order_status TINYINT, ...)) tableEnv.executeSql( CREATE TABLE enriched_orders (order_id BIGINT, order_line_id BIGINT, order_status TINYINT, ...)) val compiledPlan tableEnv.compilePlanSql( |INSERT INTO enriched_orders |SELECT a.order_id, a.order_line_id, b.order_status, ... |FROM orders a JOIN line_orders b ON a.line_order_id b.order_line_id |.stripMargin) // CompilePlan#writeToFile only supports a local file path, if you need to write to remote filesystem, // please use tableEnv.executeSql(COMPILE PLAN hdfs://path/to/plan.json FOR ...) compiledPlan.writeToFile(/path/to/plan.json)SQL CLI 方式Flink SQL CREATE TABLE orders (order_id BIGINT, order_line_id BIGINT, buyer_id BIGINT, ...); [INFO] Execute statement succeeded. Flink SQL CREATE TABLE line_orders (order_line_id BIGINT, order_status TINYINT, ...); [INFO] Execute statement succeeded. Flink SQL CREATE TABLE enriched_orders (order_id BIGINT, order_line_id BIGINT, order_status TINYINT, ...); [INFO] Execute statement succeeded. Flink SQL COMPILE PLAN file:///path/to/plan.json FOR INSERT INTO enriched_orders SELECT a.order_id, a.order_line_id, b.order_status, ... FROM orders a JOIN line_orders b ON a.order_line_id b.order_line_id; [INFO] Execute statement succeeded.COMPILE PLAN的 SQL 语法如下COMPILE PLAN [IF NOT EXISTS] plan_file_path FOR insert_statement|statement_set; statement_set: EXECUTE STATEMENT SET BEGIN insert_statement; ... insert_statement; END; insert_statement: insert_from_select|insert_from_values该语句会在指定位置生成一个 JSON 文件。除本地路径外COMPILE PLAN还支持写入hdfs://、s3://等 Flink 支持的文件系统请确保为目标写入路径设置了写入权限。第二步修改 Compiled Plan 中的状态 TTL每个状态算子会在 JSON 计划中显式生成一个名为state的数组结构如下。理论上一个拥有 k 路输入的状态算子拥有 k 个状态state: [ { index: 0, ttl: 0 ms, name: ${1st input state name} }, { index: 1, ttl: 0 ms, name: ${2nd input state name} }, ... ]找到需要修改的状态算子将 TTL 设置为带毫秒单位的正整数。例如将第一个状态算子的 TTL 设置为 1 小时{ index: 0, ttl: 3600000 ms, name: ${1st input state name} }这一 JSON 结构的字段定义可在源码flink-table/flink-table-planner/src/main/java/org/apache/flink/table/planner/plan/nodes/exec/StateMetadata.java中找到index该状态属于算子的第几路输入从 0 开始计数、ttl该路输入状态的保留时间单位毫秒、name状态描述如deduplicate-state、join-left-state等。此外该源码还实现了向后兼容逻辑若状态元数据列表为空则回退为从表配置table.exec.state.ttl读取统一 TTL。一个需要留意的经验法则下游状态算子的 TTL 不应小于上游状态算子的 TTL。第三步执行 Compiled PlanEXECUTE PLAN语句会反序列化上述 JSON 文件进一步生成 JobGraph 并提交作业。通过EXECUTE PLAN提交的作业其状态算子的 TTL 值从文件中读取配置项table.exec.state.ttl的值会被忽略。Java 方式TableEnvironment tableEnv TableEnvironment.create(EnvironmentSettings.inStreamingMode()); tableEnv.executeSql( CREATE TABLE orders (order_id BIGINT, order_line_id BIGINT, buyer_id BIGINT, ...)); tableEnv.executeSql( CREATE TABLE line_orders (order_line_id BIGINT, order_status TINYINT, ...)); tableEnv.executeSql( CREATE TABLE enriched_orders (order_id BIGINT, order_line_id BIGINT, order_status TINYINT, ...)); // PlanReference#fromFile only supports a local file path, if you need to read from remote filesystem, // please use tableEnv.executeSql(EXECUTE PLAN hdfs://path/to/plan.json).await(); tableEnv.loadPlan(PlanReference.fromFile(/path/to/plan.json)).execute().await();Scala 方式val tableEnv TableEnvironment.create(EnvironmentSettings.inStreamingMode()) tableEnv.executeSql( CREATE TABLE orders (order_id BIGINT, order_line_id BIGINT, buyer_id BIGINT, ...)) tableEnv.executeSql( CREATE TABLE line_orders (order_line_id BIGINT, order_status TINYINT, ...)) tableEnv.executeSql( CREATE TABLE enriched_orders (order_id BIGINT, order_line_id BIGINT, order_status TINYINT, ...)) // PlanReference#fromFile only supports a local file path, if you need to read from remote filesystem, // please use tableEnv.executeSql(EXECUTE PLAN hdfs://path/to/plan.json).await() tableEnv.loadPlan(PlanReference.fromFile(/path/to/plan.json)).execute().await()SQL CLI 方式Flink SQL CREATE TABLE orders (order_id BIGINT, order_line_id BIGINT, buyer_id BIGINT, ...); [INFO] Execute statement succeeded. Flink SQL CREATE TABLE line_orders (order_line_id BIGINT, order_status TINYINT, ...); [INFO] Execute statement succeeded. Flink SQL CREATE TABLE enriched_orders (order_id BIGINT, order_line_id BIGINT, order_status TINYINT, ...); [INFO] Execute statement succeeded. Flink SQL EXECUTE PLAN file:///path/to/plan.json; [INFO] Submitting SQL update statement to the cluster... [INFO] SQL update statement has been successfully submitted to the cluster: Job ID: 79fbe3fa497e4689165dd81b1d225ea8EXECUTE PLAN的 SQL 语法EXECUTE PLAN [IF EXISTS] plan_file_path;完整示例为双流 Join 的左右状态配置不同 TTL下面通过一个计算订单明细的双流 Join 作业演示完整的算子级 TTL 配置流程。① 生成 compiled plan-- left source table CREATE TABLE Orders ( order_id INT, line_order_id INT ) WITH ( connector... ); -- right source table CREATE TABLE LineOrders ( line_order_id INT, ship_mode STRING ) WITH ( connector... ); -- sink table CREATE TABLE OrdersShipInfo ( order_id INT, line_order_id INT, ship_mode STRING ) WITH ( connector ... ); COMPILE PLAN /path/to/plan.json FOR INSERT INTO OrdersShipInfo SELECT a.order_id, a.line_order_id, b.ship_mode FROM Orders a JOIN LineOrders b ON a.line_order_id b.line_order_id;生成的 JSON 文件内容如下节选关键部分{ flinkVersion : 1.18, nodes : [ { id : 1, type : stream-exec-table-source-scan_1, scanTableSource : { table : { identifier : default_catalog.default_database.Orders, resolvedTable : { ... } } }, outputType : ROWorder_id INT, line_order_id INT, description : TableSourceScan(table[[default_catalog, default_database, Orders]], fields[order_id, line_order_id]), inputProperties : [ ] }, { id : 2, type : stream-exec-exchange_1, inputProperties : [ ... ], outputType : ROWorder_id INT, line_order_id INT, description : Exchange(distribution[hash[line_order_id]]) }, { id : 3, type : stream-exec-table-source-scan_1, scanTableSource : { table : { identifier : default_catalog.default_database.LineOrders, resolvedTable : {...} } }, outputType : ROWline_order_id INT, ship_mode VARCHAR(2147483647), description : TableSourceScan(table[[default_catalog, default_database, LineOrders]], fields[line_order_id, ship_mode]), inputProperties : [ ] }, { id : 4, type : stream-exec-exchange_1, inputProperties : [ ... ], outputType : ROWline_order_id INT, ship_mode VARCHAR(2147483647), description : Exchange(distribution[hash[line_order_id]]) }, { id : 5, type : stream-exec-join_1, joinSpec : { ... }, state : [ { index : 0, ttl : 0 ms, name : leftState }, { index : 1, ttl : 0 ms, name : rightState } ], inputProperties : [ ... ], outputType : ROWorder_id INT, line_order_id INT, line_order_id0 INT, ship_mode VARCHAR(2147483647), description : Join(joinType[InnerJoin], where[(line_order_id line_order_id0)], select[order_id, line_order_id, line_order_id0, ship_mode], leftInputSpec[NoUniqueKey], rightInputSpec[NoUniqueKey]) }, { id : 6, type : stream-exec-calc_1, projection : [ ... ], condition : null, inputProperties : [ ... ], outputType : ROWorder_id INT, line_order_id INT, ship_mode VARCHAR(2147483647), description : Calc(select[order_id, line_order_id, ship_mode]) }, { id : 7, type : stream-exec-sink_1, configuration : { ... }, dynamicTableSink : { table : { identifier : default_catalog.default_database.OrdersShipInfo, resolvedTable : { ... } } }, inputChangelogMode : [ INSERT ], inputProperties : [ ... ], outputType : ROWorder_id INT, line_order_id INT, ship_mode VARCHAR(2147483647), description : Sink(table[default_catalog.default_database.OrdersShipInfo], fields[order_id, line_order_id, ship_mode]) } ], edges : [ ... ] }② 修改状态 TTL上述 JSON 中id5的 Join 算子的状态信息如下。index代表状态属于算子的第几路输入从 0 开始当前左右流的 TTL 均为0 ms表示 TTL 尚未开启state: [ { index: 0, ttl: 0 ms, name: leftState }, { index: 1, ttl: 0 ms, name: rightState } ]现在将左流 TTL 设置为3000 ms右流设置为9000 msstate: [ { index: 0, ttl: 3000 ms, name: leftState }, { index: 1, ttl: 9000 ms, name: rightState } ]③ 执行 compiled plan保存修改后使用EXECUTE PLAN语句提交作业此时提交的作业中 Join 的左右流便使用了上述不同的 TTLEXECUTE PLAN /path/to/plan.json状态化更新与演化表程序在流模式下执行时被视为标准查询它们被定义一次后将一直作为静态的端到端end-to-end管道运行。对于这种状态化管道查询语句的改动和 Flink Planner 的改动都有可能产生完全不同的执行计划这使表程序的状态化升级与演化具有挑战性。例如为了添加一个过滤谓词优化器可能决定重排 Join 或改变内部算子的 schema这会阻碍从 savepoint 的恢复——因为改变后的拓扑和算子状态的列布局与旧计划存在差异。因此查询实现者需要确保改动在优化计划前后是兼容的。可以在 SQL 中使用EXPLAIN或在 Table API 中使用table.explain()获取详情参见解释一个表由于新的优化器规则不断被添加算子变得更加高效和专用升级到更新的 Flink 版本也可能造成不兼容的计划。警告当前框架无法保证状态可以从 savepoint 映射到新的算子拓扑上。换言之savepoint 只在查询语句和 Flink 版本保持恒定的情况下才被支持。由于社区拒绝在版本补丁如1.13.1至1.13.2上对优化计划和算子拓扑进行修改的贡献将 Table API SQL 管道升级到新的 bug fix 发行版应当是安全的然而主次major-minor版本的更新如1.12至1.13不被支持。鉴于这两个限制修改查询语句、修改 Flink 版本建议在升级后、切换到实时数据之前先用历史数据对升级后的表程序做暖机即初始化验证其能否正常启动与恢复。Flink 社区正致力于通过混合源Hybrid Source让这一切换尽可能方便。延伸阅读围绕流式表程序以下文档与本文形成完整知识体系动态表动态表的核心概念是理解流式 SQL 语义的基础时间属性时间属性及其在 Table API SQL 中的使用方式时态Temporal表时态表的概念与应用流上的 Join流式场景下支持的几种 Join流上的确定性流计算确定性的解释查询配置Table API SQL 特有的全部配置项。相关源码佐证可进一步阅读flink-table/flink-table-api-java/src/main/java/org/apache/flink/table/api/config/ExecutionConfigOptions.javatable.exec.state.ttl的默认值与语义定义、flink-table/flink-table-planner/src/main/java/org/apache/flink/table/planner/hint/StateTtlHint.javaSTATE_TTL提示的解析实现、flink-table/flink-table-planner/src/main/java/org/apache/flink/table/planner/plan/nodes/exec/StateMetadata.javaCompiledPlan 中状态元数据的 JSON 结构与向后兼容逻辑以及测试flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/hints/stream/StateTtlHintTest.java。赞分享大数据流处理批处理数据工程【免费下载链接】flink项目地址https://gitcode.com/gh_mirrors/fli/flink点击查看免费下载相关推荐解决90%状态管理问题Apache Flink状态TTL配置与数据生命周期实战指南解决90%状态管理问题Apache Flink状态TTL配置与数据生命周期实战指南 你是否还在为Flink状态无限增长导致的磁盘溢出、性能下降而头疼是否因历大数据流处理批处理数据工程Apache Flink 编程模型概念透析从有状态流处理到 SQL 的四层 API 抽象Apache Flink 编程模型概念透析从有状态流处理到 SQL 的四层 API 抽象 Flink 为流式/批式处理应用程序的开发提供了从底层有状态流处理到大数据流处理批处理数据工程MVVMFramework终极指南10分钟掌握iOS优雅开发的艺术MVVMFramework终极指南10分钟掌握iOS优雅开发的艺术 想要写出优雅的iOS代码MVVMFramework正是你需要的快速开发框架这个Obje上一篇5倍性能差GLM-4推理引擎终极对决vLLM vs TensorRT-LLM技术选型指南下一篇agenix 社区贡献指南从代码提交到文档完善的完整流程创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
RELATED READING

延伸阅读

更多一线实战笔记与深度复盘,助您持续精进