You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在Flink单个作业中同时执行StatementSet和DataStream可查询流

Flink同一作业同时运行StatementSet与可查询状态流的解决方案

核心原因

  • statementSet.execute()是Table/SQL层的独立提交入口,提交时仅会打包StatementSet内部注册的SQL操作生成作业图,不会包含DataStream层定义的可查询状态逻辑
  • flinkEnv.execute()是DataStream层的提交入口,提交时仅会打包已经转换到DataStream层的算子,未被关联到流环境的StatementSet内部语句不会被加入作业拓扑

推荐解决方案(改动最小)

将原StatementSet中注册的INSERT操作转换为DataStream层的sink算子,与可查询状态逻辑统一归入StreamExecutionEnvironment的拓扑管理,最后统一调用流环境的execute方法即可同时执行两部分逻辑:

val flinkEnv = StreamExecutionEnvironment.getExecutionEnvironment();
val tableEnv = StreamTableEnvironment.create(flinkEnv);

// 替代原statementSet.addInsertSql逻辑,将插入操作转为DataStream sink
val insertSource = tableEnv.sqlQuery("SELECT * FROM INPUT_TRANSFORFM");
val outputTable = tableEnv.from("OUTPUT");
// 将待插入数据以changelog模式写入OUTPUT表,等效于原INSERT语句
tableEnv.toChangelogStream(insertSource).sinkTo(outputTable.getSink());

// 可查询状态逻辑保持不变
tableEnv.toChangelogStream(tableEnv.sqlQuery("SELECT * FROM OUTPUT"))
    .keyBy(row -> row.getField(0))
    .asQueryableState("OUTPUT_CHANGELOG_STATE");

// 统一提交,两部分逻辑都会被执行
flinkEnv.execute("job");

备选方案

如果需要保留StatementSet的使用方式,可将可查询状态逻辑封装为自定义Table Sink:

  1. 实现继承SinkFunction<Row>的自定义Sink,内部完成keyBy、可查询状态注册的逻辑
  2. 将自定义Sink注册到Table环境为临时表
  3. 在StatementSet中同时添加原INSERT语句、以及INSERT INTO 自定义可查询状态_sink SELECT * FROM OUTPUT语句
  4. 调用statementSet.execute()即可同时执行所有逻辑

内容的提问来源于stack exchange,提问作者gaurav miglani

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.09.27 16:06:00