Flink SQL批量提交DELETE语句报错及并行执行咨询
Flink批量提交DELETE语句遇到的问题
问题场景与代码实现
我尝试通过以下方式向TableEnvironment提交多条DELETE语句:
val settings = EnvironmentSettings.newInstance.inBatchMode.build() val env = TableEnvironment.create(settings) createSchema(env) val queries = getDeleteQueries queries.foreach(q => env.executeSql(q).await())
根据Flink 1.19的DELETE语句文档,这种方式应当可行,但执行第一条查询后抛出异常。
异常信息
循环提交单条语句时的异常
Caused by: org.apache.flink.util.FlinkRuntimeException: Cannot have more than one execute() or executeAsync() call in a single environment. at org.apache.flink.client.program.StreamContextEnvironment.validateAllowedExecution(StreamContextEnvironment.java:199) ~[flink-dist-1.19.1.jar:1.19.1] at org.apache.flink.client.program.StreamContextEnvironment.executeAsync(StreamContextEnvironment.java:187) ~[flink-dist-1.19.1.jar:1.19.1] at org.apache.flink.table.planner.delegation.DefaultExecutor.executeAsync(DefaultExecutor.java:110) ~[?:?] at org.apache.flink.table.api.internal.TableEnvironmentImpl.executeInternal(TableEnvironmentImpl.java:1032) ~[flink-table-api-java-uber-1.19.1.jar:1.19.1] at org.apache.flink.table.api.internal.TableEnvironmentImpl.executeInternal(TableEnvironmentImpl.java:876) ~[flink-table-api-java-uber-1.19.1.jar:1.19.1] at org.apache.flink.table.api.internal.TableEnvironmentImpl.executeInternal(TableEnvironmentImpl.java:1112) ~[flink-table-api-java-uber-1.19.1.jar:1.19.1] at org.apache.flink.table.api.internal.TableEnvironmentImpl.executeSql(TableEnvironmentImpl.java:735) ~[flink-table-api-java-uber-1.19.1.jar:1.19.1]
查看源码发现核心限制:单个环境仅支持提交一个作业,校验逻辑如下:
private void validateAllowedExecution() { if (enforceSingleJobExecution && jobCounter > 0) { throw new FlinkRuntimeException( "Cannot have more than one execute() or executeAsync() call in a single environment."); } jobCounter++; }
分号分隔多条语句提交时的异常
Caused by: java.lang.IllegalArgumentException: only single statement supported at org.apache.flink.util.Preconditions.checkArgument(Preconditions.java:138) ~[flink-dist-1.19.1.jar:1.19.1] at org.apache.flink.table.planner.delegation.ParserImpl.parse(ParserImpl.java:104) ~[?:?] at org.apache.flink.table.api.internal.TableEnvironmentImpl.executeSql(TableEnvironmentImpl.java:728) ~[flink-table-api-java-uber-1.19.1.jar:1.19.1]
环境信息
- Flink版本:1.19.1
- 部署方式:通过Flink Kubernetes Operator 1.10.0提交至Kubernetes环境
- DELETE语句示例:
DELETE FROM hdfs_table WHERE id IN ( SELECT DISTINCT a.id FROM jdbc_table_a a JOIN jdbc_table_b b ON a.b_id = b.id WHERE b.should_be_deleted )
求助问题
- 我缺失了哪些配置,导致无法提交多条DELETE语句?
- 是否可以并行提交多条DELETE语句,以复用JDBC表而非每次查询都重新加载?
内容的提问来源于stack exchange,提问作者Lyashko Kirill
相关产品推荐
相关产品推荐

