第一章.Flink SQL简介

Apache Flink 有两种关系型API 来做流批统一处理: Table API 和 SQL

Table API 是用于Scala与Java语言的查询API,它可以用一种非常直观的方式来组合使用选取,过滤,join等关系型算子

Flink SQL是基于Apache Calcite来实现的标准SQL,这两种API的查询对于批(DataSet) 和流(DataStream)的输入有相同的语义,也会产生同样的计算结果

day09[Flink SQL编程(上)] - 图1

1.动态表

  1. 动态表(Dynamic Tables)是FLink的支持流数据的Table API 和SQL的核心概念
  2. 与表示批处理数据的静态表不同,动态表是随时间变化的,可以像查询静态批处理表一样查询他们,查询动态表将生成一个连续查询(Continuous Query),一个连续查询永远不会终止,结果会生成一个动态表,查询不断更新其动态结果集以反映其(动态)输入表上的更改
  3. 需要注意的是,连续查询的结果在语义上总是等价于批处理模式在输入表快照上执行的相同查询的结果

2.连续查询

  1. 在动态表上计算一个连续查询,并生成一个新的动态表,与批处理查询不同,连续查询从不终止,并根据其输入表上的更新更新其结果表
  2. 在任何时候,连续查询的结果在语义上与以批处理模式在输入表上执行的相同查询的结果相同

day09[Flink SQL编程(上)] - 图2

  1. 当查询开始,clicks 表(左侧)是空的。
  2. 当第一行数据被插入到 clicks 表时,查询开始计算结果表。第一行数据 [Mary,./home] 插入后,结果表(右侧,上部)由一行 [Mary, 1] 组成。
  3. 当第二行 [Bob, ./cart] 插入到 clicks 表时,查询会更新结果表并插入了一行新数据 [Bob, 1]。
  4. 第三行 [Mary, ./prod?id=1] 将产生已计算的结果行的更新,[Mary, 1] 更新成 [Mary, 2]。
  5. 最后,当第四行数据加入 clicks 表时,查询将第三行 [Liz, 1] 插入到结果表中。

第二章.Flink Table API

1.导入依赖

  1. <dependency>
  2. <groupId>org.apache.flink</groupId>
  3. <artifactId>flink-table-planner-blink_${scala.binary.version}</artifactId>
  4. <version>${flink.version}</version>
  5. <scope>provided</scope>
  6. </dependency>
  7. <dependency>
  8. <groupId>org.apache.flink</groupId>
  9. <artifactId>flink-streaming-scala_${scala.binary.version}</artifactId>
  10. <version>${flink.version}</version>
  11. <scope>provided</scope>
  12. </dependency>
  13. <dependency>
  14. <groupId>org.apache.flink</groupId>
  15. <artifactId>flink-csv</artifactId>
  16. <version>${flink.version}</version>
  17. </dependency>
  18. <dependency>
  19. <groupId>org.apache.flink</groupId>
  20. <artifactId>flink-json</artifactId>
  21. <version>${flink.version}</version>
  22. </dependency>
  23. <!-- https://mvnrepository.com/artifact/org.apache.commons/commons-compress -->
  24. <dependency>
  25. <groupId>org.apache.commons</groupId>
  26. <artifactId>commons-compress</artifactId>
  27. <version>1.21</version>
  28. </dependency>

2.基本使用:表与DataStream的混合使用

写法一

  1. package com.atguigu.flink.day09;
  2. import com.atguigu.flink.day02.pojo.WaterSensor;
  3. import org.apache.flink.streaming.api.datastream.DataStream;
  4. import org.apache.flink.streaming.api.datastream.DataStreamSource;
  5. import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
  6. import org.apache.flink.table.api.Table;
  7. import org.apache.flink.table.api.bridge.java.StreamTableEnvironment;
  8. import org.apache.flink.types.Row;
  9. public class $01_TableBaseUse {
  10. public static void main(String[] args) throws Exception {
  11. //获取流的执行环境
  12. StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
  13. env.setParallelism(1);
  14. //读取集合中的数据
  15. DataStreamSource<WaterSensor> stream = env.fromElements(
  16. new WaterSensor("sensor_1", 1000L, 10),
  17. new WaterSensor("sensor_1", 2000L, 20),
  18. new WaterSensor("sensor_2", 3000L, 30),
  19. new WaterSensor("sensor_1", 4000L, 40),
  20. new WaterSensor("sensor_1", 5000L, 50),
  21. new WaterSensor("sensor_2", 6000L, 60)
  22. );
  23. //获取表的执行环境
  24. StreamTableEnvironment tenv = StreamTableEnvironment.create(env);
  25. //把流转成动态表
  26. Table table = tenv.fromDataStream(stream);
  27. table.printSchema();//打印表的元数据
  28. //在动态表上进行连续查询
  29. Table result = table.where("id=='sensor_1'")
  30. .select("id,vc");
  31. //把查询结果(动态表)转成流
  32. DataStream<Row> resultStream = tenv.toAppendStream(result, Row.class);
  33. //把流进行sink
  34. resultStream.print();
  35. env.execute();
  36. }
  37. }

写法二(常用)

  1. package com.atguigu.flink.day09;
  2. import com.atguigu.flink.day02.pojo.WaterSensor;
  3. import org.apache.flink.streaming.api.datastream.DataStream;
  4. import org.apache.flink.streaming.api.datastream.DataStreamSource;
  5. import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
  6. import org.apache.flink.table.api.Table;
  7. import org.apache.flink.table.api.bridge.java.StreamTableEnvironment;
  8. import org.apache.flink.types.Row;
  9. import static org.apache.flink.table.api.Expressions.$;
  10. public class $02_TableBaseUseNormal {
  11. public static void main(String[] args) throws Exception {
  12. //获取流的执行环境
  13. StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
  14. env.setParallelism(1);
  15. //读取集合中的数据
  16. DataStreamSource<WaterSensor> stream = env.fromElements(
  17. new WaterSensor("sensor_1", 1000L, 10),
  18. new WaterSensor("sensor_1", 2000L, 20),
  19. new WaterSensor("sensor_2", 3000L, 30),
  20. new WaterSensor("sensor_1", 4000L, 40),
  21. new WaterSensor("sensor_1", 5000L, 50),
  22. new WaterSensor("sensor_2", 6000L, 60)
  23. );
  24. //获取表的执行环境
  25. StreamTableEnvironment tenv = StreamTableEnvironment.create(env);
  26. //把流转成动态表
  27. Table table = tenv.fromDataStream(stream);
  28. //在动态表上进行连续查询
  29. Table result = table.where($("id").isEqual("sensor_1"))
  30. .select($("id").as("id1"), $("vc"));
  31. result.execute().print();
  32. }
  33. }

3.基本使用:聚合操作

  1. package com.atguigu.flink.day09;
  2. import com.atguigu.flink.day02.pojo.WaterSensor;
  3. import org.apache.flink.streaming.api.datastream.DataStreamSource;
  4. import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
  5. import org.apache.flink.table.api.AggregatedTable;
  6. import org.apache.flink.table.api.Table;
  7. import org.apache.flink.table.api.bridge.java.StreamTableEnvironment;
  8. import static org.apache.flink.table.api.Expressions.$;
  9. public class $03_TableBaseUseAgg {
  10. public static void main(String[] args) throws Exception {
  11. //获取流的执行环境
  12. StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
  13. env.setParallelism(1);
  14. //读取集合中的数据
  15. DataStreamSource<WaterSensor> stream = env.fromElements(
  16. new WaterSensor("sensor_1", 1000L, 10),
  17. new WaterSensor("sensor_1", 2000L, 20),
  18. new WaterSensor("sensor_2", 3000L, 30),
  19. new WaterSensor("sensor_1", 4000L, 40),
  20. new WaterSensor("sensor_1", 5000L, 50),
  21. new WaterSensor("sensor_2", 6000L, 60)
  22. );
  23. //获取表的执行环境
  24. StreamTableEnvironment tenv = StreamTableEnvironment.create(env);
  25. //把流转成动态表
  26. Table table = tenv.fromDataStream(stream);
  27. //在动态表上进行连续查询
  28. /*Table result = table.groupBy($("id"))
  29. .aggregate($("vc").sum().as("vc_sum"))
  30. .select($("id"), $("vc_sum"));*/
  31. Table result = table.groupBy($("id"))
  32. .select($("id"), $("vc").sum().as("vc_sum"));
  33. result.execute().print();
  34. /**
  35. * 把动态表转换成流. 如果涉及到数据的更新, 要用到撤回流. 多个了一个boolean标记
  36. * DataStream<Tuple2<Boolean, Row>> resultStream = tenv.toRetractStream(result, Row.class);
  37. * resultStream.print()
  38. */
  39. }
  40. }

4.表到流的转换

  1. 动态表可以像普通数据库表一样通过insert,update,和delete来不断修改,它可能是一个只有一行,不断更新的表,也可能是一个insert-only的表,没有update和delete修改,或者介于两者之间的其他表
  2. 在将动态表转换为流,或将其写入外部系统时,需要对这些更改进行编码,Flink的Table API和SQL支持三种方式来编码一个动态表的变化
  3. Append-only流:仅通过insert操作修改的动态表可以通过输出输入的行转换为流
  4. Retract流:retract流包含两种类型的message:addmessages 和 retract messages,通过将insert操作编码为addmessage,将delete操作编码为retract message,将update操作编码为更新行的retract message和更新行的add message,将动态表转换为retract流
  5. Upsert流:upsert流包含两种类型的message:upsertmessages和delete messages.转换为upsert流的动态表需要(可能是组合的)唯一键,通过将insert和upsert操作操作编码为upsert message,将delete操作编码为delete message,将具有唯一键的动态表转换为流,消费流的算子需要知道唯一键的属性,以便正确的应用message.与retract流的主要区别在于update操作是用单个message编码的,因此效率更高

5.使用Table API读写文件

  1. package com.atguigu.flink.day09;
  2. import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
  3. import org.apache.flink.table.api.DataTypes;
  4. import org.apache.flink.table.api.Table;
  5. import org.apache.flink.table.api.bridge.java.StreamTableEnvironment;
  6. import org.apache.flink.table.descriptors.Csv;
  7. import org.apache.flink.table.descriptors.FileSystem;
  8. import org.apache.flink.table.descriptors.Schema;
  9. import static org.apache.flink.table.api.Expressions.$;
  10. public class $04_TableToFile {
  11. public static void main(String[] args) throws Exception {
  12. //获取流的执行环境
  13. StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
  14. env.setParallelism(1);
  15. //获取表的执行环境
  16. StreamTableEnvironment tenv = StreamTableEnvironment.create(env);
  17. //表的元数据信息
  18. Schema schema = new Schema()
  19. .field("id", DataTypes.STRING())
  20. .field("ts", DataTypes.BIGINT())
  21. .field("vc", DataTypes.INT());
  22. //通过读取文件创建临时表
  23. tenv.connect(new FileSystem().path("input/sensor.txt"))
  24. //文件的切分方式
  25. .withFormat(new Csv().lineDelimiter("\n").fieldDelimiter(','))
  26. .withSchema(schema)
  27. .createTemporaryTable("sensor");
  28. Table resultTable = tenv.from("sensor")
  29. .where($("id").isEqual("sensor_1"))
  30. .select($("id"), $("vc"));
  31. //建立动态表与文件(sink)进行关联
  32. tenv.connect(new FileSystem().path("input/a.txt"))
  33. //文件的切分方式
  34. .withFormat(new Csv().fieldDelimiter(','))
  35. .withSchema(
  36. new Schema()
  37. .field("id", DataTypes.STRING())
  38. .field("vc", DataTypes.INT())
  39. )
  40. .createTemporaryTable("a");
  41. //把数据写入到与sink关联的临时表中(动态表)
  42. resultTable.executeInsert("a");
  43. }
  44. }

6.使用Table API读写Kafka

  1. package com.atguigu.flink.day09;
  2. import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
  3. import org.apache.flink.table.api.DataTypes;
  4. import org.apache.flink.table.api.Table;
  5. import org.apache.flink.table.api.bridge.java.StreamTableEnvironment;
  6. import org.apache.flink.table.descriptors.*;
  7. import static org.apache.flink.table.api.Expressions.$;
  8. public class $05_TableToKafka {
  9. public static void main(String[] args) throws Exception {
  10. //获取流的执行环境
  11. StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
  12. env.setParallelism(1);
  13. //获取表的执行环境
  14. StreamTableEnvironment tenv = StreamTableEnvironment.create(env);
  15. //表的元数据信息
  16. Schema schema = new Schema()
  17. .field("id", DataTypes.STRING())
  18. .field("ts", DataTypes.BIGINT())
  19. .field("vc", DataTypes.INT());
  20. //通过读取kafka创建临时表
  21. tenv.connect(
  22. new Kafka()
  23. .version("universal")
  24. .property("bootstrap.servers","hadoop162:9092")
  25. .property("group.id","$05_TableToKafka")
  26. .topic("s1")
  27. .startFromLatest()
  28. )
  29. //文件的切分方式
  30. .withFormat(new Json())
  31. .withSchema(schema)
  32. .createTemporaryTable("s1");
  33. Table s1 = tenv.from("s1").select($("id"), $("vc"));
  34. //与sink关联
  35. tenv.connect(
  36. new Kafka()
  37. .version("universal")
  38. .property("bootstrap.servers","hadoop162:9092")
  39. .topic("s2")
  40. .sinkPartitionerRoundRobin()
  41. )
  42. .withFormat(new Csv().lineDelimiter(""))
  43. .withSchema(
  44. new Schema()
  45. .field("id", DataTypes.STRING())
  46. .field("vc", DataTypes.INT())
  47. ).createTemporaryTable("s2");
  48. s1.executeInsert("s2");
  49. }
  50. }

第三章.Flink SQL

1.基本使用

  1. package com.atguigu.flink.day09;
  2. import com.atguigu.flink.day02.pojo.WaterSensor;
  3. import org.apache.flink.streaming.api.datastream.DataStream;
  4. import org.apache.flink.streaming.api.datastream.DataStreamSource;
  5. import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
  6. import org.apache.flink.table.api.Table;
  7. import org.apache.flink.table.api.bridge.java.StreamTableEnvironment;
  8. import org.apache.flink.types.Row;
  9. public class $06_SQLBaseUse {
  10. public static void main(String[] args) throws Exception {
  11. //获取流的执行环境
  12. StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
  13. env.setParallelism(1);
  14. //读取集合中的数据
  15. DataStreamSource<WaterSensor> stream = env.fromElements(
  16. new WaterSensor("sensor_1", 1000L, 10),
  17. new WaterSensor("sensor_1", 2000L, 20),
  18. new WaterSensor("sensor_2", 3000L, 30),
  19. new WaterSensor("sensor_1", 4000L, 40),
  20. new WaterSensor("sensor_1", 5000L, 50),
  21. new WaterSensor("sensor_2", 6000L, 60)
  22. );
  23. //获取表的执行环境
  24. StreamTableEnvironment tenv = StreamTableEnvironment.create(env);
  25. //把流转成动态表
  26. Table table = tenv.fromDataStream(stream);
  27. //查询未注册的表
  28. /**
  29. * tenv.executeSql("");执行DDL和DML(insert,update和delete)
  30. */
  31. //注意from和where之后要有空格
  32. Table t1 = tenv.sqlQuery("select * from " + table + " where id= 'sensor_1'");//只执行查询
  33. t1.execute().print();
  34. //查询已注册的表
  35. //1.先为table对象注册一个临时表
  36. tenv.createTemporaryView("sensor",table);
  37. //2.执行SQL查询
  38. tenv.sqlQuery("select * from sensor where id= 'sensor_1'").execute().print();
  39. }
  40. }

2.使用SQL API读写文件

  1. package com.atguigu.flink.day09;
  2. import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
  3. import org.apache.flink.table.api.bridge.java.StreamTableEnvironment;
  4. public class $07_SQLToFile {
  5. public static void main(String[] args) throws Exception {
  6. //获取流的执行环境
  7. StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
  8. env.setParallelism(1);
  9. //获取表的执行环境
  10. StreamTableEnvironment tenv = StreamTableEnvironment.create(env);
  11. //通过ddl建立一个动态 表,与file source 进行关联
  12. tenv.executeSql("create table sensor(id string, ts bigint, vc int" +")with(" +
  13. " 'connector' = 'filesystem', "+
  14. " 'path' = 'input/sensor.txt', "+
  15. " 'format'= 'csv' " +
  16. ")"
  17. );
  18. //建立动态表与sink进行关联
  19. tenv.executeSql("create table b(id string, vc int" +")with(" +
  20. " 'connector' = 'filesystem', "+
  21. " 'path' = 'input/b.txt', "+
  22. " 'format'= 'csv' " +
  23. ")"
  24. );
  25. tenv.executeSql("insert into b select id, vc from sensor where id='sensor_1' ");
  26. }
  27. }

3.使用SQL API读写Kafka

  1. package com.atguigu.flink.day09;
  2. import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
  3. import org.apache.flink.table.api.Table;
  4. import org.apache.flink.table.api.bridge.java.StreamTableEnvironment;
  5. public class $08_SQLToKafka {
  6. public static void main(String[] args) throws Exception {
  7. //获取流的执行环境
  8. StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
  9. env.setParallelism(1);
  10. //获取表的执行环境
  11. StreamTableEnvironment tenv = StreamTableEnvironment.create(env);
  12. tenv.executeSql("create table sensor(id string, ts bigint, vc int" +")with(" +
  13. " 'connector' = 'kafka',\n "+
  14. " 'properties.bootstrap.servers' = 'hadoop162:9092',\n "+
  15. " 'properties.group.id'= '$08_SQLToKafka',\n " +
  16. " 'topic' = 's1',\n"+
  17. " 'scan.startup.mode' = 'latest-offset',\n"+
  18. " 'format' = 'csv'" +
  19. ")"
  20. );
  21. tenv.executeSql("create table ss(id string, vc int" +")with(" +
  22. " 'connector' = 'kafka',\n "+
  23. " 'properties.bootstrap.servers' = 'hadoop162:9092',\n "+
  24. " 'topic' = 's2',\n"+
  25. " 'format' = 'json',\n" +
  26. " 'sink.partitioner' = 'round-robin' "+
  27. ")"
  28. );
  29. Table result = tenv.sqlQuery("select id, vc from sensor");
  30. result.executeInsert("ss");
  31. }
  32. }

第四章.时间属性

1.处理时间

DataStream到Table转换时定义

  1. package com.atguigu.flink.day09;
  2. import com.atguigu.flink.day02.pojo.WaterSensor;
  3. import org.apache.flink.streaming.api.datastream.DataStreamSource;
  4. import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
  5. import org.apache.flink.table.api.Table;
  6. import org.apache.flink.table.api.bridge.java.StreamTableEnvironment;
  7. import static org.apache.flink.table.api.Expressions.$;
  8. public class $09_TimeProcessing {
  9. public static void main(String[] args) throws Exception {
  10. //获取流的执行环境
  11. StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
  12. env.setParallelism(1);
  13. //读取集合中的数据
  14. DataStreamSource<WaterSensor> stream = env.fromElements(
  15. new WaterSensor("sensor_1", 1000L, 10),
  16. new WaterSensor("sensor_1", 2000L, 20),
  17. new WaterSensor("sensor_2", 3000L, 30),
  18. new WaterSensor("sensor_1", 4000L, 40),
  19. new WaterSensor("sensor_1", 5000L, 50),
  20. new WaterSensor("sensor_2", 6000L, 60)
  21. );
  22. //获取表的执行环境
  23. StreamTableEnvironment tenv = StreamTableEnvironment.create(env);
  24. //必须添加一个新的字段作为处理时间
  25. Table table = tenv.fromDataStream(stream, $("id"), $("ts"), $("vc"), $("pt").proctime());
  26. table.execute().print();
  27. }
  28. }

在创建表的DDL中定义

  1. package com.atguigu.flink.day09;
  2. import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
  3. import org.apache.flink.table.api.bridge.java.StreamTableEnvironment;
  4. public class $10_TimeDDLProcessing {
  5. public static void main(String[] args) throws Exception {
  6. //获取流的执行环境
  7. StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
  8. env.setParallelism(1);
  9. //获取表的执行环境
  10. StreamTableEnvironment tenv = StreamTableEnvironment.create(env);
  11. tenv.executeSql(
  12. "create table sensor(" +
  13. " id string, " +
  14. " ts bigint, " +
  15. " vc int, " +
  16. " pt as proctime() " +
  17. ")with(" +
  18. " 'connector' = 'filesystem', " +
  19. " 'path' = 'input/sensor.txt', " +
  20. " 'format' = 'csv' " +
  21. ")"
  22. );
  23. tenv.sqlQuery("select * from sensor").execute().print();
  24. }
  25. }

2.事件时间

DataStream到Table转换时定义

  1. package com.atguigu.flink.day09;
  2. import com.atguigu.flink.day02.pojo.WaterSensor;
  3. import org.apache.flink.api.common.eventtime.WatermarkStrategy;
  4. import org.apache.flink.streaming.api.datastream.SingleOutputStreamOperator;
  5. import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
  6. import org.apache.flink.table.api.Table;
  7. import org.apache.flink.table.api.bridge.java.StreamTableEnvironment;
  8. import java.time.Duration;
  9. import static org.apache.flink.table.api.Expressions.$;
  10. public class $11_TimeEventTime {
  11. public static void main(String[] args) throws Exception {
  12. //获取流的执行环境
  13. StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
  14. env.setParallelism(1);
  15. //读取集合中的数据
  16. SingleOutputStreamOperator<WaterSensor> stream = env.fromElements(
  17. new WaterSensor("sensor_1", 1000L, 10),
  18. new WaterSensor("sensor_1", 2000L, 20),
  19. new WaterSensor("sensor_2", 3000L, 30),
  20. new WaterSensor("sensor_1", 4000L, 40),
  21. new WaterSensor("sensor_1", 5000L, 50),
  22. new WaterSensor("sensor_2", 6000L, 60)
  23. ).assignTimestampsAndWatermarks(
  24. WatermarkStrategy.<WaterSensor>forBoundedOutOfOrderness(Duration.ofSeconds(3))
  25. .withTimestampAssigner((ws, ts) -> ws.getTs())
  26. );
  27. //获取表的执行环境
  28. StreamTableEnvironment tenv = StreamTableEnvironment.create(env);
  29. //直接在原来属性指定为事件时间
  30. Table table = tenv.fromDataStream(stream, $("id"), $("ts").rowtime(), $("vc"));
  31. table.execute().print();
  32. }
  33. }

在创建表的DDL中定义

  1. package com.atguigu.flink.day09;
  2. import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
  3. import org.apache.flink.table.api.bridge.java.StreamTableEnvironment;
  4. import java.time.ZoneOffset;
  5. public class $12_TimeDDLEventTime {
  6. public static void main(String[] args) throws Exception {
  7. //获取流的执行环境
  8. StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
  9. env.setParallelism(1);
  10. //获取表的执行环境
  11. StreamTableEnvironment tenv = StreamTableEnvironment.create(env);
  12. //设置时区
  13. tenv.getConfig().setLocalTimeZone(ZoneOffset.ofHours(0));
  14. tenv.executeSql(
  15. "create table sensor(" +
  16. " id string, " +
  17. " ts bigint, " +
  18. " vc int, " +
  19. " et as to_timestamp(from_unixtime(ts/1000)), " +
  20. " WATERMARK FOR et AS et - INTERVAL '3' SECOND " +
  21. ")with(" +
  22. " 'connector' = 'filesystem', " +
  23. " 'path' = 'input/sensor.txt', " +
  24. " 'format' = 'csv' " +
  25. ")"
  26. );
  27. tenv.sqlQuery("select * from sensor").execute().print();
  28. }
  29. }

第五章.窗口

1.Table API中使用窗口

一.Group Windows

  1. package com.atguigu.flink.day09;
  2. import com.atguigu.flink.day02.pojo.WaterSensor;
  3. import org.apache.flink.api.common.eventtime.WatermarkStrategy;
  4. import org.apache.flink.streaming.api.datastream.SingleOutputStreamOperator;
  5. import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
  6. import org.apache.flink.table.api.*;
  7. import org.apache.flink.table.api.bridge.java.StreamTableEnvironment;
  8. import java.time.Duration;
  9. import static org.apache.flink.table.api.Expressions.$;
  10. import static org.apache.flink.table.api.Expressions.lit;
  11. public class $13_TableWindowGroup {
  12. public static void main(String[] args) throws Exception {
  13. //获取流的执行环境
  14. StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
  15. env.setParallelism(1);
  16. //读取集合中的数据
  17. SingleOutputStreamOperator<WaterSensor> stream = env.fromElements(
  18. new WaterSensor("sensor_1", 1000L, 10),
  19. new WaterSensor("sensor_1", 2000L, 20),
  20. new WaterSensor("sensor_2", 3000L, 30),
  21. new WaterSensor("sensor_1", 4000L, 40),
  22. new WaterSensor("sensor_1", 5000L, 50),
  23. new WaterSensor("sensor_2", 6000L, 60)
  24. ).assignTimestampsAndWatermarks(
  25. WatermarkStrategy.<WaterSensor>forBoundedOutOfOrderness(Duration.ofSeconds(3))
  26. .withTimestampAssigner((ws, ts) -> ws.getTs())
  27. );
  28. //获取表的执行环境
  29. StreamTableEnvironment tenv = StreamTableEnvironment.create(env);
  30. //直接在原来属性指定为事件时间
  31. Table table = tenv.fromDataStream(stream, $("id"), $("ts").rowtime(), $("vc"));
  32. //滚动窗口
  33. //GroupWindow window = Tumble.over(lit(3).second()).on($("ts")).as("w");
  34. //滑动窗口
  35. GroupWindow window = Slide.over(lit(5).second()).every(lit(2).second()).on($("ts")).as("w");
  36. //会话窗口
  37. //GroupWindow window = Session.withGap(lit(2).second()).on($("ts")).as("w");
  38. table
  39. .window(window)
  40. .groupBy($("id"),$("w"))
  41. .select($("id"),$("w").start(),$("w").end(),$("vc").sum().as("sum_vc"))
  42. .execute()
  43. .print();
  44. }
  45. }

二.Over Windows

  1. package com.atguigu.flink.day09;
  2. import com.atguigu.flink.day02.pojo.WaterSensor;
  3. import org.apache.flink.api.common.eventtime.WatermarkStrategy;
  4. import org.apache.flink.streaming.api.datastream.SingleOutputStreamOperator;
  5. import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
  6. import org.apache.flink.table.api.*;
  7. import org.apache.flink.table.api.bridge.java.StreamTableEnvironment;
  8. import org.apache.flink.table.expressions.Expression;
  9. import org.apache.flink.table.runtime.operators.over.frame.OverWindowFrame;
  10. import java.time.Duration;
  11. import static org.apache.flink.table.api.Expressions.*;
  12. public class $14_TableWindowOver {
  13. public static void main(String[] args) throws Exception {
  14. //获取流的执行环境
  15. StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
  16. env.setParallelism(1);
  17. //读取集合中的数据
  18. SingleOutputStreamOperator<WaterSensor> stream = env.fromElements(
  19. new WaterSensor("sensor_1", 1000L, 10),
  20. new WaterSensor("sensor_1", 2000L, 20),
  21. new WaterSensor("sensor_2", 3000L, 30),
  22. new WaterSensor("sensor_1", 4000L, 40),
  23. new WaterSensor("sensor_1", 4000L, 50),
  24. new WaterSensor("sensor_2", 6000L, 60)
  25. ).assignTimestampsAndWatermarks(
  26. WatermarkStrategy.<WaterSensor>forBoundedOutOfOrderness(Duration.ofSeconds(3))
  27. .withTimestampAssigner((ws, ts) -> ws.getTs())
  28. );
  29. //获取表的执行环境
  30. StreamTableEnvironment tenv = StreamTableEnvironment.create(env);
  31. //直接在原来属性指定为事件时间
  32. Table table = tenv.fromDataStream(stream, $("id"), $("ts").rowtime(), $("vc"));
  33. //基于行数开窗
  34. OverWindow over = Over.partitionBy($("id")).orderBy($("ts")).preceding(Expressions.UNBOUNDED_ROW).as("w");
  35. //基于时间开窗
  36. OverWindow over1 = Over.partitionBy($("id")).orderBy($("ts")).preceding(Expressions.UNBOUNDED_RANGE).as("w1");
  37. //当事件时间向前算1s得到一个窗口
  38. OverWindow over2 = Over.partitionBy($("id")).orderBy($("ts")).preceding(lit(1).second()).as("w3");
  39. //当前行向前推算一行算一个窗口
  40. OverWindow over3 = Over.partitionBy($("id")).orderBy($("ts")).preceding(rowInterval(1L)).as("w4");
  41. table.window(over1)
  42. .select($("id"),$("ts"),$("vc").sum().over($("w1")).as("vc_sum"))
  43. .execute()
  44. .print();
  45. }
  46. }

2.SQL API中使用窗口

一.Group Windows

  1. package com.atguigu.flink.day09;
  2. import com.atguigu.flink.day02.pojo.WaterSensor;
  3. import org.apache.flink.api.common.eventtime.WatermarkStrategy;
  4. import org.apache.flink.streaming.api.datastream.SingleOutputStreamOperator;
  5. import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
  6. import org.apache.flink.table.api.Table;
  7. import org.apache.flink.table.api.bridge.java.StreamTableEnvironment;
  8. import java.time.Duration;
  9. import static org.apache.flink.table.api.Expressions.$;
  10. public class $15_SQLWindowGroup {
  11. public static void main(String[] args) throws Exception {
  12. //获取流的执行环境
  13. StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
  14. env.setParallelism(1);
  15. //读取集合中的数据
  16. SingleOutputStreamOperator<WaterSensor> stream = env.fromElements(
  17. new WaterSensor("sensor_1", 1000L, 10),
  18. new WaterSensor("sensor_1", 2000L, 20),
  19. new WaterSensor("sensor_2", 3000L, 30),
  20. new WaterSensor("sensor_1", 4000L, 40),
  21. new WaterSensor("sensor_1", 5000L, 50),
  22. new WaterSensor("sensor_2", 6000L, 60)
  23. ).assignTimestampsAndWatermarks(
  24. WatermarkStrategy.<WaterSensor>forBoundedOutOfOrderness(Duration.ofSeconds(3))
  25. .withTimestampAssigner((ws, ts) -> ws.getTs())
  26. );
  27. //获取表的执行环境
  28. StreamTableEnvironment tenv = StreamTableEnvironment.create(env);
  29. //直接在原来属性指定为事件时间
  30. Table table = tenv.fromDataStream(stream, $("id"), $("ts").rowtime(), $("vc"));
  31. tenv.createTemporaryView("sensor", table);
  32. // 滚动窗口
  33. /*tenv.sqlQuery("select " +
  34. " id," +
  35. " tumble_start(ts, interval '3' second) stt, " +
  36. " tumble_end(ts, interval '3' second) edt, " +
  37. " sum(vc) " +
  38. "from sensor " +
  39. "group by id, TUMBLE(ts, interval '3' second)").execute().print();*/
  40. //跳跃的时间窗口(滑动窗口)
  41. /*tenv.sqlQuery("select " +
  42. " id," +
  43. " hop_start(ts, interval '3' second, interval '5' second) stt, " +
  44. " hop_end(ts, interval '3' second, interval '5' second) edt, " +
  45. " sum(vc) " +
  46. "from sensor " +
  47. "group by id, HOP(ts, interval '3' second, interval '5' second)").execute().print();*/
  48. //会话的时间窗口
  49. tenv.sqlQuery("select " +
  50. " id," +
  51. " session_start(ts, interval '2' second) stt, " +
  52. " session_end(ts, interval '2' second) edt, " +
  53. " sum(vc) " +
  54. "from sensor " +
  55. "group by id, SESSION(ts, interval '2' second)").execute().print();
  56. }
  57. }

二.Over Windows

  1. package com.atguigu.flink.day09;
  2. import com.atguigu.flink.day02.pojo.WaterSensor;
  3. import org.apache.flink.api.common.eventtime.WatermarkStrategy;
  4. import org.apache.flink.streaming.api.datastream.SingleOutputStreamOperator;
  5. import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
  6. import org.apache.flink.table.api.Table;
  7. import org.apache.flink.table.api.bridge.java.StreamTableEnvironment;
  8. import java.time.Duration;
  9. import static org.apache.flink.table.api.Expressions.*;
  10. public class $16_SQLWindowOver {
  11. public static void main(String[] args) throws Exception {
  12. //获取流的执行环境
  13. StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
  14. env.setParallelism(1);
  15. //读取集合中的数据
  16. SingleOutputStreamOperator<WaterSensor> stream = env.fromElements(
  17. new WaterSensor("sensor_1", 1000L, 10),
  18. new WaterSensor("sensor_1", 2000L, 20),
  19. new WaterSensor("sensor_2", 3000L, 30),
  20. new WaterSensor("sensor_1", 4000L, 40),
  21. new WaterSensor("sensor_1", 4000L, 50),
  22. new WaterSensor("sensor_2", 6000L, 60)
  23. ).assignTimestampsAndWatermarks(
  24. WatermarkStrategy.<WaterSensor>forBoundedOutOfOrderness(Duration.ofSeconds(3))
  25. .withTimestampAssigner((ws, ts) -> ws.getTs())
  26. );
  27. //获取表的执行环境
  28. StreamTableEnvironment tenv = StreamTableEnvironment.create(env);
  29. //直接在原来属性指定为事件时间
  30. Table table = tenv.fromDataStream(stream, $("id"), $("ts").rowtime(), $("vc"));
  31. tenv.createTemporaryView("sensor", table);
  32. /*tenv.sqlQuery("select" +
  33. " id, " +
  34. " ts," +
  35. " vc," +
  36. // " sum(vc) over(partition by id order by ts rows between unbounded preceding and current row) vc_sum " +
  37. // " sum(vc) over(partition by id order by ts rows between 1 preceding and current row) vc_sum " +
  38. // " sum(vc) over(partition by id order by ts range between unbounded preceding and current row) vc_sum " +
  39. // " sum(vc) over(partition by id order by ts range between interval '1' second preceding and current row) vc_sum " +
  40. " sum(vc) over(partition by id order by ts) vc_sum " +
  41. "from sensor")
  42. .execute()
  43. .print();*/
  44. tenv.sqlQuery("select" +
  45. " id, " +
  46. " ts," +
  47. " vc," +
  48. " sum(vc) over w vc_sum, " +
  49. " max(vc) over w vc_max, " +
  50. " min(vc) over w vc_min " +
  51. "from sensor " +
  52. "window w as(partition by id order by ts rows between unbounded preceding and current row)")
  53. .execute()
  54. .print();
  55. }
  56. }