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

如何在本地执行模式下运行Flink Table API(Java+IDEA运行)

1. 依赖配置(Maven为例)

在你的pom.xml中添加以下依赖,确保所有Flink组件版本一致(示例使用1.18.0稳定版,可根据需求调整):

<dependencies>
    <!-- Flink Table API 核心依赖 -->
    <dependency>
        <groupId>org.apache.flink</groupId>
        <artifactId>flink-table-api-java</artifactId>
        <version>1.18.0</version>
        <scope>compile</scope>
    </dependency>
    <!-- Table Planner:负责SQL/Table API的执行计划生成 -->
    <dependency>
        <groupId>org.apache.flink</groupId>
        <artifactId>flink-table-planner-loader</artifactId>
        <version>1.18.0</version>
        <scope>runtime</scope>
    </dependency>
    <!-- 本地执行所需的运行时依赖 -->
    <dependency>
        <groupId>org.apache.flink</groupId>
        <artifactId>flink-runtime</artifactId>
        <version>1.18.0</version>
        <scope>runtime</scope>
    </dependency>
    <!-- 流处理场景依赖(若使用批处理可忽略) -->
    <dependency>
        <groupId>org.apache.flink</groupId>
        <artifactId>flink-streaming-java</artifactId>
        <version>1.18.0</version>
        <scope>runtime</scope>
    </dependency>
    <!-- 测试用连接器:可选datagen(生成模拟数据)或filesystem(读取本地文件) -->
    <dependency>
        <groupId>org.apache.flink</groupId>
        <artifactId>flink-connector-datagen</artifactId>
        <version>1.18.0</version>
        <scope>runtime</scope>
    </dependency>
</dependencies>

2. 本地执行的Java代码示例

以下是一个完整的本地执行示例,无需依赖任何Flink集群组件,所有计算在当前JVM进程内完成:

import org.apache.flink.table.api.EnvironmentSettings;
import org.apache.flink.table.api.Table;
import org.apache.flink.table.api.TableEnvironment;
import static org.apache.flink.table.api.Expressions.$;

public class LocalFlinkTableDemo {
    public static void main(String[] args) {
        // 配置本地执行环境:可选批处理(inBatchMode)或流处理(inStreamingMode)
        EnvironmentSettings settings = EnvironmentSettings
                .newInstance()
                .inBatchMode()
                .useLocalExecution() // 明确指定本地执行模式(默认本地运行即为该模式)
                .build();

        // 创建TableEnvironment
        TableEnvironment tableEnv = TableEnvironment.create(settings);

        // 注册模拟数据源表(无需外部文件,自动生成测试数据)
        tableEnv.executeSql("CREATE TABLE user_data (" +
                "user_id INT," +
                "user_name STRING," +
                "age INT" +
                ") WITH (" +
                "'connector' = 'datagen'," +
                "'rows-per-second' = '5'," +
                "'fields.user_id.kind' = 'sequence'," +
                "'fields.user_id.start' = '1'," +
                "'fields.user_id.end' = '20'," +
                "'fields.age.min' = '16'," +
                "'fields.age.max' = '40'" +
                ")");

        // 执行Table API查询:筛选年龄≥18的用户
        Table adultUsers = tableEnv.from("user_data")
                .where($("age").isGreaterOrEqual(18))
                .select($("user_id"), $("user_name"), $("age"));

        // 将查询结果打印到控制台
        adultUsers.execute().print();
    }
}

3. 在IntelliJ中运行程序

  1. 等待Maven/Gradle完成依赖下载,确保项目无编译错误
  2. 找到LocalFlinkTableDemo类的main方法,右键选择Run 'LocalFlinkTableDemo.main()'
  3. 程序会直接在IntelliJ控制台输出查询结果,全程无需启动JobManager、TaskManager等集群组件

关键说明

  • 本地执行模式下,所有计算逻辑在当前JVM进程内运行,资源由JVM自行管理,适合开发、测试阶段快速验证逻辑
  • 若需切换流处理模式,只需将inBatchMode()替换为inStreamingMode(),流处理场景下程序会持续运行并输出结果
  • 避免在生产环境使用本地模式,大规模数据处理需部署Flink集群

内容的提问来源于stack exchange,提问作者Aly Ayman

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.03 16:05:16