如何在本地执行模式下运行Flink Table API(Java+IDEA运行)
在本地执行模式下运行Flink Table API(Java + IntelliJ)
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中运行程序
- 等待Maven/Gradle完成依赖下载,确保项目无编译错误
- 找到
LocalFlinkTableDemo类的main方法,右键选择Run 'LocalFlinkTableDemo.main()' - 程序会直接在IntelliJ控制台输出查询结果,全程无需启动JobManager、TaskManager等集群组件
关键说明
- 本地执行模式下,所有计算逻辑在当前JVM进程内运行,资源由JVM自行管理,适合开发、测试阶段快速验证逻辑
- 若需切换流处理模式,只需将
inBatchMode()替换为inStreamingMode(),流处理场景下程序会持续运行并输出结果 - 避免在生产环境使用本地模式,大规模数据处理需部署Flink集群
内容的提问来源于stack exchange,提问作者Aly Ayman
相关产品推荐
相关产品推荐

