如何使用Flink本地执行模式执行批量SQL?
Flink本地CollectionEnvironment执行批量SQL操作步骤
前置依赖配置
首先在项目中引入对应版本的Flink依赖(以1.14版本、Java语言、Scala 2.12版本为例),Maven配置参考如下:
<dependencies> <dependency> <groupId>org.apache.flink</groupId> <artifactId>flink-java</artifactId> <version>1.14.6</version> </dependency> <dependency> <groupId>org.apache.flink</groupId> <artifactId>flink-table-api-java-bridge_2.12</artifactId> <version>1.14.6</version> </dependency> <dependency> <groupId>org.apache.flink</groupId> <artifactId>flink-table-planner_2.12</artifactId> <version>1.14.6</version> </dependency> <dependency> <groupId>org.apache.flink</groupId> <artifactId>flink-clients_2.12</artifactId> <version>1.14.6</version> </dependency> </dependencies>
核心操作流程
- 初始化本地集合执行环境,以及对应批量SQL执行环境:
// 初始化CollectionEnvironment本地执行环境 ExecutionEnvironment batchEnv = ExecutionEnvironment.createCollectionsEnvironment(); // 对接批量Table环境 BatchTableEnvironment tableEnv = BatchTableEnvironment.create(batchEnv);
- 注册数据源,你可以把本地的测试集合转成Flink表,也可以读取本地测试文件作为输入,示例如下:
// 构造本地测试数据集 List<Row> testData = Arrays.asList( Row.of(1, "张三", 25), Row.of(2, "李四", 30), Row.of(3, "王五", 28) ); // 定义字段类型 TypeInformation<Row> rowType = Types.ROW_NAMED( new String[]{"id", "name", "age"}, Types.INT, Types.STRING, Types.INT ); // 转成DataSet并注册为临时视图 DataSet<Row> sourceDs = batchEnv.fromCollection(testData, rowType); tableEnv.createTemporaryView("user_info", sourceDs);
- 执行批量SQL逻辑,和线上保持一致即可:
// 执行SQL查询,示例:统计年龄大于25的用户数量 Table resultTable = tableEnv.sqlQuery("SELECT count(1) as user_cnt FROM user_info WHERE age > 25"); // 把结果转成DataSet输出 DataSet<Row> resultDs = tableEnv.toDataSet(resultTable, Row.class); resultDs.print();
- 触发任务执行:
batchEnv.execute("本地批量SQL执行任务");
注意事项
- 所有Flink相关依赖版本必须统一,避免版本兼容问题
- 如果线上使用了自定义UDF、自定义序列化规则、特殊的SQL配置(比如时区、并行度、数值精度规则),本地环境也要做完全一致的配置,保证两边逻辑执行结果一致
- CollectionEnvironment仅适合本地小批量测试使用,不要用于生产环境的大量数据计算
内容的提问来源于stack exchange,提问作者zuston
相关产品推荐
相关产品推荐

