一下的狀態(tài)管理、狀態(tài) TTL 與算子級生命周期配置)
大數(shù)據(jù)流處理批處理數(shù)據(jù)工程【免費下載鏈接】flink項目地址https://gitcode.com/gh_mirrors/fli/flink點擊查看免費下載Flink 的 Table API 與 SQL 是一套流批統(tǒng)一的聲明式 API在有限的批式輸入和無限的流式輸入下具備相同的語義。由于關(guān)系代數(shù)與 SQL 最初是為批處理設(shè)計的流式場景下的狀態(tài)管理成為理解與調(diào)優(yōu)的關(guān)鍵。本文以docs/content.zh/docs/dev/table/concepts/overview.md為核心系統(tǒng)講解流式表程序的狀態(tài)使用方式、空閑狀態(tài)維持時間State TTL的三種配置途徑、從 Flink v1.18 起支持的算子級狀態(tài) TTL含 CompiledPlan 完整實操并延伸狀態(tài)化更新與演化等進階話題幫助你掌握狀態(tài)維度下的流式 SQL 生產(chǎn)實踐。流批統(tǒng)一狀態(tài)是流式表程序的靈魂Flink 的 Table API 與 SQL 在批式與流式輸入下共享同一套語義。區(qū)別在于批式查詢天然擁有有限的輸入集可以一次性完成計算而流式查詢面對的是無限的數(shù)據(jù)流必須以**連續(xù)查詢Continuous Query**的形式持續(xù)運行因此必須依賴狀態(tài)來保存跨時間維度的中間結(jié)果。一個流模式下運行的表程序Table program可以完整利用 Flink 作為有狀態(tài)流處理器的能力配置不同的 state backend如 RocksDB、Heap以適配不同規(guī)模的狀態(tài)存儲需求配置多種 checkpoint 選項以滿足不同的容錯與恢復(fù)需求對正在運行的 Table API SQL 管道生成 savepoint并在之后用其恢復(fù)應(yīng)用狀態(tài)。狀態(tài)使用聲明式管道中的隱式狀態(tài)由于 Table API SQL 程序是聲明式的狀態(tài)會在哪里、如何被使用并不直接可見。**Planner優(yōu)化器**負責(zé)判斷是否需要狀態(tài)來得到正確的計算結(jié)果并盡可能把管道優(yōu)化成使用更少狀態(tài)的形式。從概念上講源表從來不會在狀態(tài)中被完全保存——實現(xiàn)者處理的是邏輯表即動態(tài)表Dynamic Table各算子的狀態(tài)完全取決于具體用到的操作。顯式狀態(tài)算子Join、聚合與去重包含連接Join、聚合Aggregation或去重Deduplication等操作的語句需要在 Flink 抽象的容錯存儲內(nèi)保存中間結(jié)果這類算子被稱為狀態(tài)算子。例如對兩個表執(zhí)行普通 Join基于正確的 SQL 語義運行時假設(shè)兩表會在任意時間點進行匹配因此算子需要保存兩個表的全部輸入。為了控制狀態(tài)規(guī)模Flink 提供了優(yōu)化窗口 Join 和時段 Join利用 watermarks 概念即時間屬性讓過期的數(shù)據(jù)不再參與匹配從而顯著縮小狀態(tài)。另一個經(jīng)典例子是詞頻統(tǒng)計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是分組的鍵連續(xù)查詢?yōu)槊總€觀察到的word維護一個中間狀態(tài)來保存當(dāng)前詞頻。由于輸入word的值隨時間變化且查詢持續(xù)運行Flink 會為每個word維護一個中間狀態(tài)總狀態(tài)量會隨著新word的出現(xiàn)不斷增長——這正是流式聚合狀態(tài)下最需要警惕的內(nèi)存與存儲風(fēng)險點。隱式狀態(tài)算子SELECT 也可能引入狀態(tài)形如SELECT ... FROM ... WHERE這種只包含字段映射或過濾器的查詢通常是無狀態(tài)的。但在某些情況下根據(jù)輸入數(shù)據(jù)的特征或配置狀態(tài)算子會被隱式地推導(dǎo)出來輸入表是不帶UPDATE_BEFORE的更新流詳見表到流的轉(zhuǎn)換或配置了table-exec-source-cdc-events-duplicate。下面的例子展示了對 upsert-kafka 源表執(zhí)行最簡單的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。因此雖然查詢本身不包含任何狀態(tài)計算優(yōu)化器依然會隱式地推導(dǎo)出一個 ChangelogNormalize 狀態(tài)算子來生成完整的 changelog??臻e狀態(tài)維持時間table.exec.state.ttl空閑狀態(tài)維持時間參數(shù)table.exec.state.ttl定義了狀態(tài)的鍵在被更新后要保持多長時間才被移除。在上述詞頻例子中某個word的計數(shù)會在配置的時間內(nèi)未更新時被立刻移除。該配置項的源碼定義位于flink-table/flink-table-api-java/src/main/java/org/apache/flink/table/api/config/ExecutionConfigOptions.java鍵名為table.exec.state.ttl類型為 Duration默認值為0 ms含義是永不清理狀態(tài)清理狀態(tài)會引入額外的簿記開銷bookkeeping因此默認關(guān)閉。移除狀態(tài)的鍵之后連續(xù)查詢會完全忘記它曾經(jīng)見過這個鍵如果一條記錄帶有一個曾被移除狀態(tài)的鍵該記錄會被當(dāng)作對應(yīng)鍵的第一條記錄處理。在詞頻例子中這意味著cnt會再次從0開始計數(shù)——這是配置 TTL 時最容易忽略的語義影響。指定狀態(tài)生命周期的三種方式從 Flink v1.18 開始Table API SQL 支持多種粒度的狀態(tài) TTL 配置方式下表源自原文檔總結(jié)了它們的適用面與優(yōu)先級配置方式TableAPI/SQL 支持生效范圍優(yōu)先級SET table.exec.state.ttl ...TableAPI、SQL作業(yè)粒度默認情況下所有狀態(tài)算子都會使用該值控制狀態(tài)生命周期默認配置可被覆蓋SELECT /* STATE_TTL(...) */ ...SQL有限算子粒度當(dāng)前支持連接和分組聚合算子該值優(yōu)先作用于相應(yīng)算子的狀態(tài)生命周期詳見狀態(tài)生命周期提示修改序列化為 JSON 的 CompiledPlanTableAPI、SQL通用算子粒度可修改任一狀態(tài)算子的生命周期table.exec.state.ttl與STATE_TTL的值會序列化到 CompiledPlan若作業(yè)使用 CompiledPlan 提交最終生效的生命周期由最后一次修改的狀態(tài)元數(shù)據(jù)決定STATE_TTL 查詢提示STATE_TTL提示以 SQL 注釋形式作用于具體算子。從源碼flink-table/flink-table-planner/src/main/java/org/apache/flink/table/planner/hint/StateTtlHint.java可以看出其實現(xiàn)要點對于雙輸入算子如 Join支持形如STATE_TTL(T1 1d, T2 2d)的鍵值對寫法分別指定左右輸入的 TTL源碼中LEFT_INPUT映射為輸入側(cè) 0其余映射為輸入側(cè) 1對于單輸入算子如分組聚合支持STATE_TTL(T1 2d)形式TTL 值支持d天、h小時等 Flink 時間單位通過TimeUtils.parseDuration解析為毫秒。對應(yīng)的測試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.等情況。配置算子粒度的狀態(tài) TTL高級特性注意這是一個需要小心使用的高級特性。它僅適用于作業(yè)中使用了多個狀態(tài)、且每個狀態(tài)需要不同 TTL 的場景。無狀態(tài)作業(yè)無需關(guān)注若作業(yè)僅使用一個狀態(tài)僅需設(shè)置作業(yè)級 TTL 參數(shù)table.exec.state.ttl即可。從 Flink v1.18 開始Table API SQL 支持以每個狀態(tài)算子的入邊數(shù)為粒度配置細粒度狀態(tài) TTLOneInputStreamOperator單輸入可配置一個狀態(tài)的 TTLTwoInputStreamOperator如雙流 Join可分別為左狀態(tài)和右狀態(tài)配置 TTL更一般地具有 K 個輸入的MultipleInputStreamOperator可以配置 K 個狀態(tài) TTL。典型使用場景為雙流 Join的左右流配置不同 TTL雙流 Join 會生成擁有兩條輸入邊的TwoInputStreamOperator狀態(tài)算子分別用兩個狀態(tài)保存來自左流和右流的更新在同一作業(yè)中為不同的狀態(tài)計算設(shè)置不同 TTL例如一個 ETL 作業(yè)先用ROW_NUMBER進行去重再用GROUP BY進行聚合會生成兩個擁有單條輸入邊的OneInputStreamOperator狀態(tài)算子可為它們分別設(shè)置不同的 TTL。需要說明的是基于窗口的操作如窗口連接、窗口聚合、窗口 Top-N 等和 Interval Join 不依賴table.exec.state.ttl控制狀態(tài)保留因此它們的狀態(tài)無法在算子級別配置。第一步生成 Compiled Plan配置過程首先使用COMPILE PLAN語句生成一個 JSON 文件它表示序列化后的執(zhí)行計劃。注意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 支持的文件系統(tǒng)請確保為目標寫入路徑設(shè)置了寫入權(quán)限。第二步修改 Compiled Plan 中的狀態(tài) TTL每個狀態(tài)算子會在 JSON 計劃中顯式生成一個名為state的數(shù)組結(jié)構(gòu)如下。理論上一個擁有 k 路輸入的狀態(tài)算子擁有 k 個狀態(tài)state: [ { index: 0, ttl: 0 ms, name: ${1st input state name} }, { index: 1, ttl: 0 ms, name: ${2nd input state name} }, ... ]找到需要修改的狀態(tài)算子將 TTL 設(shè)置為帶毫秒單位的正整數(shù)。例如將第一個狀態(tài)算子的 TTL 設(shè)置為 1 小時{ index: 0, ttl: 3600000 ms, name: ${1st input state name} }這一 JSON 結(jié)構(gòu)的字段定義可在源碼flink-table/flink-table-planner/src/main/java/org/apache/flink/table/planner/plan/nodes/exec/StateMetadata.java中找到index該狀態(tài)屬于算子的第幾路輸入從 0 開始計數(shù)、ttl該路輸入狀態(tài)的保留時間單位毫秒、name狀態(tài)描述如deduplicate-state、join-left-state等。此外該源碼還實現(xiàn)了向后兼容邏輯若狀態(tài)元數(shù)據(jù)列表為空則回退為從表配置table.exec.state.ttl讀取統(tǒng)一 TTL。一個需要留意的經(jīng)驗法則下游狀態(tài)算子的 TTL 不應(yīng)小于上游狀態(tài)算子的 TTL。第三步執(zhí)行 Compiled PlanEXECUTE PLAN語句會反序列化上述 JSON 文件進一步生成 JobGraph 并提交作業(yè)。通過EXECUTE PLAN提交的作業(yè)其狀態(tài)算子的 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 的左右狀態(tài)配置不同 TTL下面通過一個計算訂單明細的雙流 Join 作業(yè)演示完整的算子級 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 文件內(nèi)容如下節(jié)選關(guān)鍵部分{ 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 : [ ... ] }② 修改狀態(tài) TTL上述 JSON 中id5的 Join 算子的狀態(tài)信息如下。index代表狀態(tài)屬于算子的第幾路輸入從 0 開始當(dāng)前左右流的 TTL 均為0 ms表示 TTL 尚未開啟state: [ { index: 0, ttl: 0 ms, name: leftState }, { index: 1, ttl: 0 ms, name: rightState } ]現(xiàn)在將左流 TTL 設(shè)置為3000 ms右流設(shè)置為9000 msstate: [ { index: 0, ttl: 3000 ms, name: leftState }, { index: 1, ttl: 9000 ms, name: rightState } ]③ 執(zhí)行 compiled plan保存修改后使用EXECUTE PLAN語句提交作業(yè)此時提交的作業(yè)中 Join 的左右流便使用了上述不同的 TTLEXECUTE PLAN /path/to/plan.json狀態(tài)化更新與演化表程序在流模式下執(zhí)行時被視為標準查詢它們被定義一次后將一直作為靜態(tài)的端到端end-to-end管道運行。對于這種狀態(tài)化管道查詢語句的改動和 Flink Planner 的改動都有可能產(chǎn)生完全不同的執(zhí)行計劃這使表程序的狀態(tài)化升級與演化具有挑戰(zhàn)性。例如為了添加一個過濾謂詞優(yōu)化器可能決定重排 Join 或改變內(nèi)部算子的 schema這會阻礙從 savepoint 的恢復(fù)——因為改變后的拓撲和算子狀態(tài)的列布局與舊計劃存在差異。因此查詢實現(xiàn)者需要確保改動在優(yōu)化計劃前后是兼容的??梢栽?SQL 中使用EXPLAIN或在 Table API 中使用table.explain()獲取詳情參見解釋一個表由于新的優(yōu)化器規(guī)則不斷被添加算子變得更加高效和專用升級到更新的 Flink 版本也可能造成不兼容的計劃。警告當(dāng)前框架無法保證狀態(tài)可以從 savepoint 映射到新的算子拓撲上。換言之savepoint 只在查詢語句和 Flink 版本保持恒定的情況下才被支持。由于社區(qū)拒絕在版本補丁如1.13.1至1.13.2上對優(yōu)化計劃和算子拓撲進行修改的貢獻將 Table API SQL 管道升級到新的 bug fix 發(fā)行版應(yīng)當(dāng)是安全的然而主次major-minor版本的更新如1.12至1.13不被支持。鑒于這兩個限制修改查詢語句、修改 Flink 版本建議在升級后、切換到實時數(shù)據(jù)之前先用歷史數(shù)據(jù)對升級后的表程序做暖機即初始化驗證其能否正常啟動與恢復(fù)。Flink 社區(qū)正致力于通過混合源Hybrid Source讓這一切換盡可能方便。延伸閱讀圍繞流式表程序以下文檔與本文形成完整知識體系動態(tài)表動態(tài)表的核心概念是理解流式 SQL 語義的基礎(chǔ)時間屬性時間屬性及其在 Table API SQL 中的使用方式時態(tài)Temporal表時態(tài)表的概念與應(yīng)用流上的 Join流式場景下支持的幾種 Join流上的確定性流計算確定性的解釋查詢配置Table API SQL 特有的全部配置項。相關(guān)源碼佐證可進一步閱讀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提示的解析實現(xiàn)、flink-table/flink-table-planner/src/main/java/org/apache/flink/table/planner/plan/nodes/exec/StateMetadata.javaCompiledPlan 中狀態(tài)元數(shù)據(jù)的 JSON 結(jié)構(gòu)與向后兼容邏輯以及測試flink-table/flink-table-planner/src/test/java/org/apache/flink/table/planner/plan/hints/stream/StateTtlHintTest.java。贊分享大數(shù)據(jù)流處理批處理數(shù)據(jù)工程【免費下載鏈接】flink項目地址https://gitcode.com/gh_mirrors/fli/flink點擊查看免費下載相關(guān)推薦解決90%狀態(tài)管理問題Apache Flink狀態(tài)TTL配置與數(shù)據(jù)生命周期實戰(zhàn)指南解決90%狀態(tài)管理問題Apache Flink狀態(tài)TTL配置與數(shù)據(jù)生命周期實戰(zhàn)指南 你是否還在為Flink狀態(tài)無限增長導(dǎo)致的磁盤溢出、性能下降而頭疼是否因歷大數(shù)據(jù)流處理批處理數(shù)據(jù)工程Apache Flink 編程模型概念透析從有狀態(tài)流處理到 SQL 的四層 API 抽象Apache Flink 編程模型概念透析從有狀態(tài)流處理到 SQL 的四層 API 抽象 Flink 為流式/批式處理應(yīng)用程序的開發(fā)提供了從底層有狀態(tài)流處理到大數(shù)據(jù)流處理批處理數(shù)據(jù)工程MVVMFramework終極指南10分鐘掌握iOS優(yōu)雅開發(fā)的藝術(shù)MVVMFramework終極指南10分鐘掌握iOS優(yōu)雅開發(fā)的藝術(shù) 想要寫出優(yōu)雅的iOS代碼MVVMFramework正是你需要的快速開發(fā)框架這個Obje上一篇5倍性能差GLM-4推理引擎終極對決vLLM vs TensorRT-LLM技術(shù)選型指南下一篇agenix 社區(qū)貢獻指南從代碼提交到文檔完善的完整流程創(chuàng)作聲明:本文部分內(nèi)容由AI輔助生成(AIGC),僅供參考