Flink 1.19.2本地运行Table API报错:无法实例化执行器
问题
尝试搭建本地Flink任务,通过Table API从Kafka数据源获取数据并打印输出,代码、错误信息及依赖配置如下:
任务代码
public static void main(String[] args) { // Standard Flink streaming execution environment StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(); // EnvironmentSettings with Blink planner in STREAMING mode EnvironmentSettings settings = EnvironmentSettings.newInstance() .inStreamingMode() // Explicitly use streaming mode .build(); // Blink planner is default in recent Flink versions // Create StreamTableEnvironment with explicit settings StreamTableEnvironment sTableEnv = StreamTableEnvironment.create(env, settings); Schema tableSchema = Schema.newBuilder() .column("id", DataTypes.STRING()) .column("value", DataTypes.DOUBLE()) .build(); TableDescriptor descriptor = TableDescriptor.forConnector("kafka") // Connector type .schema(tableSchema) .format("json") // Format of the data .option("topic", "my-topic") .option("properties.bootstrap.servers", "localhost:9092") .option("scan.startup.mode", "earliest-offset") .build(); sTableEnv.createTemporaryTable("TestTable", descriptor); // Define Print connector as a sink table TableDescriptor printSink = TableDescriptor .forConnector("print") .schema(tableSchema) .build(); sTableEnv.createTemporaryTable("PrintSink", printSink); // Read from Kafka table Table resultTable = sTableEnv.from("TestTable"); // Write the result to the print sink resultTable.executeInsert("PrintSink"); }
错误信息
Exception in thread "main" org.apache.flink.table.api.TableException: Could not instantiate the executor. Make sure a planner module is on the classpath at org.apache.flink.table.api.bridge.internal.AbstractStreamTableEnvironmentImpl.lookupExecutor(AbstractStreamTableEnvironmentImpl.java:109) at org.apache.flink.table.api.bridge.java.internal.StreamTableEnvironmentImpl.create(StreamTableEnvironmentImpl.java:110) at org.apache.flink.table.api.bridge.java.StreamTableEnvironment.create(StreamTableEnvironment.java:122) at com.ruoethren.flink.poc.Main.main(Main.java:22) Caused by: org.apache.flink.table.api.ValidationException: Could not find any factories that implement 'org.apache.flink.table.delegation.ExecutorFactory' in the classpath. at org.apache.flink.table.factories.FactoryUtil.discoverFactory(FactoryUtil.java:605) at org.apache.flink.table.api.bridge.internal.AbstractStreamTableEnvironmentImpl.lookupExecutor(AbstractStreamTableEnvironmentImpl.java:106) ... 3 more
导入包
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment; import org.apache.flink.table.api.DataTypes; import org.apache.flink.table.api.EnvironmentSettings; import org.apache.flink.table.api.Schema; import org.apache.flink.table.api.Table; import org.apache.flink.table.api.TableDescriptor; import org.apache.flink.table.api.bridge.java.StreamTableEnvironment;
POM依赖
<dependencies> <!-- https://mvnrepository.com/artifact/org.apache.flink/flink-java --> <dependency> <groupId>org.apache.flink</groupId> <artifactId>flink-streaming-java</artifactId> <version>${flink.version}</version> </dependency> <dependency> <groupId>org.apache.flink</groupId> <artifactId>flink-table-api-java-bridge</artifactId> <version>${flink.version}</version> </dependency> <dependency> <groupId>org.apache.flink</groupId> <artifactId>flink-table-runtime</artifactId> <version>${flink.version}</version> </dependency> <dependency> <groupId>org.apache.flink</groupId> <artifactId>flink-connector-kafka</artifactId> <version>3.3.0-1.19</version> </dependency> <dependency> <groupId>org.apache.flink</groupId> <artifactId>flink-connector-files</artifactId> <version>${flink.version}</version> <scope>provided</scope> </dependency> <dependency> <groupId>org.apache.flink</groupId> <artifactId>flink-json</artifactId> <version>${flink.version}</version> </dependency> <dependency> <groupId>org.apache.flink</groupId> <artifactId>flink-clients</artifactId> <version>${flink.version}</version> </dependency> </dependencies>
解决方案
错误核心是类路径中缺少Table API的Planner实现依赖,同时需保证依赖版本一致性,具体修复步骤:
添加Blink Planner加载器依赖:在POM中加入以下依赖,版本必须与
${flink.version}保持一致。Flink 1.15+版本推荐使用该依赖自动处理Planner类加载:<dependency> <groupId>org.apache.flink</groupId> <artifactId>flink-table-planner-loader</artifactId> <version>${flink.version}</version> </dependency>统一Kafka Connector版本:当前Kafka Connector版本
3.3.0-1.19需与Flink主版本匹配(例如Flink 1.19.x对应此版本),若${flink.version}为其他版本,需调整为兼容的Connector版本,避免版本冲突。检查依赖范围:本地调试时,若
flink-connector-files是任务必需的,可移除其<scope>provided</scope>配置,确保依赖能被正确加载。
调整后重新构建项目并运行,即可解决该错误。
内容的提问来源于stack exchange,提问作者Ruoethren Pugunisparam
相关产品推荐
相关产品推荐

