如何在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:
- 实现继承
SinkFunction<Row>的自定义Sink,内部完成keyBy、可查询状态注册的逻辑 - 将自定义Sink注册到Table环境为临时表
- 在StatementSet中同时添加原INSERT语句、以及
INSERT INTO 自定义可查询状态_sink SELECT * FROM OUTPUT语句 - 调用
statementSet.execute()即可同时执行所有逻辑
内容的提问来源于stack exchange,提问作者gaurav miglani
相关产品推荐
相关产品推荐

