flink1.12.x Table API 和 SQL:创建工程,读取csv文件,创建表,实现查询(connector连接器) 作者:马育民 • 2021-10-16 20:36 • 阅读:10891 # 说明 文中案例,要读取 csv 文件,然后通过 SQL 查询该文件中的数据 # 创建maven工程 略 # maven 设置maven # pom.xml ### 关键依赖 在之前的依赖基础上,增加下面依赖: ``` <!-- table api、sql需要的依赖 --> <dependency> <groupId>org.apache.flink</groupId> <artifactId>flink-table-planner-blink_2.12</artifactId> <version>${flink.version}</version> </dependency> <!-- 读取csv文件,需要下面依赖 --> <dependency> <groupId>org.apache.flink</groupId> <artifactId>flink-csv</artifactId> <version>${flink.version}</version> </dependency> ``` ### 完整配置 ``` <properties> <project.build.sourceEncoding>UTF-8</project.build.sourceEncoding> <maven.compiler.source>1.8</maven.compiler.source> <maven.compiler.target>1.8</maven.compiler.target> <flink.version>1.12.0</flink.version> </properties> <dependencies> <!--java--> <dependency> <groupId>org.apache.flink</groupId> <artifactId>flink-java</artifactId> <version>${flink.version}</version> </dependency> <dependency> <groupId>org.apache.flink</groupId> <artifactId>flink-streaming-java_2.12</artifactId> <version>${flink.version}</version> </dependency> <dependency> <groupId>org.apache.flink</groupId> <artifactId>flink-clients_2.12</artifactId> <version>${flink.version}</version> </dependency> <dependency> <groupId>org.slf4j</groupId> <artifactId>slf4j-api</artifactId> <version>1.7.30</version> </dependency> <dependency> <groupId>org.slf4j</groupId> <artifactId>slf4j-log4j12</artifactId> <version>1.7.30</version> </dependency> <dependency> <groupId>log4j</groupId> <artifactId>log4j</artifactId> <version>1.2.17</version> </dependency> <!-- kafka Dependency --> <dependency> <groupId>org.apache.flink</groupId> <artifactId>flink-connector-kafka_2.12</artifactId> <version>${flink.version}</version> </dependency> <dependency> <groupId>org.apache.flink</groupId> <artifactId>flink-connector-files</artifactId> <version>${flink.version}</version> </dependency> <dependency> <groupId>org.apache.flink</groupId> <artifactId>flink-runtime-web_2.12</artifactId> <version>${flink.version}</version> </dependency> <dependency> <groupId>org.apache.flink</groupId> <artifactId>flink-table-planner-blink_2.12</artifactId> <version>${flink.version}</version> </dependency> <dependency> <groupId>org.apache.flink</groupId> <artifactId>flink-csv</artifactId> <version>${flink.version}</version> </dependency> </dependencies> ``` # log4j.properties **提示:**不加会有警告提示 在 `resources` 目录下创建 `log4j.properties`,内容如下: ``` # DEBUG 级别,即:所有信息都会输出 log4j.rootLogger=WARN, stdout # 配置 stdout,输出到控制台 log4j.appender.stdout=org.apache.log4j.ConsoleAppender log4j.appender.stdout.layout=org.apache.log4j.PatternLayout log4j.appender.stdout.layout.ConversionPattern=%d %p [%c] - %m%n ``` # 准备数据 在工程下创建文件夹 `data`,然后创建文件 `1.csv`,内容如下: ``` java从入门到精通,50,李雷 java从入门到精通,80,李雷 hadoop从入门到精通,90,韩梅梅 flink从入门到精通,70,韩梅梅 mysql从入门到精通,100,李雷 hive从入门到精通,100,lucy html从入门到精通,120,lucy ``` # java ``` package test; import org.apache.flink.api.common.RuntimeExecutionMode; import org.apache.flink.api.java.tuple.Tuple2; import org.apache.flink.streaming.api.datastream.DataStream; import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment; import org.apache.flink.table.api.EnvironmentSettings; import org.apache.flink.table.api.Table; import org.apache.flink.table.api.TableResult; import org.apache.flink.table.api.bridge.java.StreamTableEnvironment; import org.apache.flink.types.Row; public class TestTable读取文件 { public static void main(String[] args) throws Exception { StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(); env.setRuntimeMode(RuntimeExecutionMode.AUTOMATIC);//自动模式,根据数据源自行判断 //构建Table环境 EnvironmentSettings bsSettings = EnvironmentSettings.newInstance().useBlinkPlanner().inStreamingMode().build(); StreamTableEnvironment stEnv = StreamTableEnvironment.create(env, bsSettings); /* 创建表 注意:表名不要出现关键字,如:order等,否则会报错 要读取 csv文件,需要依赖flink-csv */ TableResult table=stEnv.executeSql("CREATE TABLE t_order (" + " name STRING," + " price double," + " username STRING" + ") WITH ( " + " 'connector' = 'filesystem'," + // 连接器,表示文件系统 " 'path' = 'file:///D:/bigdata/flink_table_sql/data'," + // 路径指定到文件夹 或 文件都可以 " 'format' = 'csv'," + // 文件格式 " 'csv.field-delimiter' = ','" + // csv列之间间隔符号 ")"); // 查询price大于60的数据,返回 table 对象 Table res=stEnv.sqlQuery("select * from t_order where price>60"); System.out.println("Schema信息:"); //打印Schema res.printSchema(); //将计算后的数据,append到ds DataStream<Row> dsRes2 = stEnv.toAppendStream(res, Row.class); dsRes2.print("---"); env.execute(); } } ``` ### 注意: - 建表时不要出现关键字,如:`order` 等 # 执行结果 ``` Schema信息: root |-- name: STRING |-- price: DOUBLE |-- username: STRING ---:4> flink从入门到精通,70.0,韩梅梅 ---:1> hadoop从入门到精通,90.0,韩梅梅 ---:5> java从入门到精通,80.0,李雷 ---:8> mysql从入门到精通,100.0,李雷 ---:3> html从入门到精通,120.0,lucy ---:6> hive从入门到精通,100.0,lucy ``` 原文出处:/show_1IX23NIaLLkI.html