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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.13 01:44:57