第一章.Catalog

  1. catalog提供了元数据信息,例如数据库,表,分区,视图以及数据库或其他外部系统中存储的函数和信息
  2. 数据处理最关键的方面之一是管理元数据,元数据可以是临时的,例如临时表,或者通过TableEnvironment注册的UDF,元数据也可以是持久化的,例如HiveMetaStore中的元数据,Catalog提供了一个统一的API,用于管理元数据,并使其可以从TableAPI和SQL查询语句中来访问

1.Catalog类型

  1. GenericInMemoryCatalog
  2. GenericInMemoryCatalog 是基于内存实现的 Catalog,所有元数据只在 session 的生命周期内可用。
  3. JdbcCatalog
  4. JdbcCatalog 使得用户可以将 Flink 通过 JDBC 协议连接到关系数据库。PostgresCatalog 是当前实现的唯一一种 JDBC Catalog。
  5. HiveCatalog
  6. HiveCatalog 有两个用途:作为原生 Flink 元数据的持久化存储,以及作为读写现有 Hive 元数据的接口。 Flink 的 Hive 文档 提供了有关设置 HiveCatalog 以及访问现有 Hive 元数据的详细信息。

2.HiveCatalog

一.导入依赖

  1. <dependency>
  2. <groupId>org.apache.flink</groupId>
  3. <artifactId>flink-connector-hive_${scala.binary.version}</artifactId>
  4. <version>${flink.version}</version>
  5. </dependency>
  6. <!-- Hive Dependency -->
  7. <dependency>
  8. <groupId>org.apache.hive</groupId>
  9. <artifactId>hive-exec</artifactId>
  10. <version>3.1.2</version>
  11. </dependency>

二.在hadoop162上启动元数据服务

  1. nohup hive --service metastore >/dev/null 2>&1 &

三.代码编写

  1. package com.atguigu.flink.day10;
  2. import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
  3. import org.apache.flink.table.api.bridge.java.StreamTableEnvironment;
  4. import org.apache.flink.table.catalog.hive.HiveCatalog;
  5. public class $01_Hive {
  6. public static void main(String[] args) {
  7. StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
  8. env.setParallelism(1);
  9. StreamTableEnvironment tenv = StreamTableEnvironment.create(env);
  10. //1.创建一个hive catalog
  11. HiveCatalog hive = new HiveCatalog("hive", "default", "input/");
  12. //2.注册hive catalog
  13. tenv.registerCatalog("hive",hive);
  14. //3.设置默认的catalog
  15. tenv.useCatalog("hive");
  16. tenv.useDatabase("default");
  17. tenv.sqlQuery("select * from student").execute().print();
  18. }
  19. }

第二章.函数

Flink 允许用户在Table API 和 SQL中使用函数来进行数据的转换

1.内置函数

Flink Table API 和SQL给用户提供了大量的函数用于数据转换

具体参照官网: https://ci.apache.org/projects/flink/flink-docs-release-1.12/dev/table/functions/systemFunctions.html

2.自定义函数

  1. 自定义函数(UDF)是一种扩展开发机制,可以用来查询语句中调用难以用其他方式表达的频繁使用和自定义的逻辑
  2. 自定义函数分类:
  3. 1.标量函数:标量值转换成一个新标量值
  4. 2.表值函数:将标量值转换成新的行数据
  5. 3.聚合函数:将多行数据里的标量值转换成一个新的标量值
  6. 4.表值聚合函数:将多行数据里的标量值转换成新的行数据
  7. 5.异步表值函数:是异步查询外部数据系统的特殊函数
  8. 函数用于SQL查询前要先经过注册,而在用于Table API时,函数可以先注册后调用,也可以内联后直接使用

一.标量函数

  1. 介绍:
  2. 用户定义的标量函数,可以将0、1或多个标量值,映射到新的标量值。
  3. 为了定义标量函数,必须在org.apache.flink.table.functions中扩展基类Scalar Function,并实现(一个或多个)求值(evaluation,eval)方法。标量函数的行为由求值方法决定,求值方法必须公开声明并命名为eval(直接def声明,没有override)。求值方法的参数类型和返回类型,确定了标量函数的参数和返回类型。
  1. package com.atguigu.flink.day10;
  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.DataStreamSource;
  5. import org.apache.flink.streaming.api.datastream.SingleOutputStreamOperator;
  6. import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
  7. import org.apache.flink.table.api.Table;
  8. import org.apache.flink.table.api.bridge.java.StreamTableEnvironment;
  9. import org.apache.flink.table.functions.ScalarFunction;
  10. import static org.apache.flink.table.api.Expressions.$;
  11. import static org.apache.flink.table.api.Expressions.call;
  12. /**
  13. * 变成大写字母的标量函数
  14. */
  15. public class $02_FunctionScalar {
  16. public static void main(String[] args) throws Exception {
  17. //获取流的执行环境
  18. StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
  19. env.setParallelism(1);
  20. //读取集合中的数据
  21. DataStreamSource<WaterSensor> stream = env.fromElements(
  22. new WaterSensor("sensor_1", 1000L, 10),
  23. new WaterSensor("sensor_1", 2000L, 20),
  24. new WaterSensor("sensor_2", 3000L, 30),
  25. new WaterSensor("sensor_1", 4000L, 40),
  26. new WaterSensor("sensor_1", 4000L, 50),
  27. new WaterSensor("sensor_2", 6000L, 60)
  28. );
  29. //获取表的执行环境
  30. StreamTableEnvironment tenv = StreamTableEnvironment.create(env);
  31. Table table = tenv.fromDataStream(stream);
  32. //1.在table api中使用
  33. //1.1内联的方式
  34. /*table
  35. .select($("id"),call(MyUpperCase.class,$("id")).as("id_upper"))
  36. .execute()
  37. .print();*/
  38. //1.2注册后使用
  39. /*tenv.createTemporaryFunction("toUpper",MyUpperCase.class);
  40. table
  41. .select($("id"),call("toUpper",$("id")).as("id_upper"))
  42. .execute()
  43. .print();*/
  44. //2.在SQL语句中使用
  45. //2.1先注册
  46. tenv.createTemporaryFunction("toUpper",MyUpperCase.class);
  47. //2.2再使用
  48. tenv.sqlQuery(" select id, toUpper(id) from " + table).execute().print();
  49. }
  50. public static class MyUpperCase extends ScalarFunction{
  51. public String eval(String s){
  52. return s==null?null:s.toUpperCase();
  53. }
  54. }
  55. }

二.表值函数

  1. 跟自定义标量函数一样,自定义表值函数的输入参数也可以是 0 到多个标量。但是跟标量函数只能返回一个值不同的是,它可以返回任意多行。返回的每一行可以包含 1 到多列,如果输出行只包含 1 列,会省略结构化信息并生成标量值,这个标量值在运行阶段会隐式地包装进行里。
  2. 要定义一个表值函数,你需要扩展 org.apache.flink.table.functions 下的 TableFunction,可以通过实现多个名为 eval 的方法对求值方法进行重载。像其他函数一样,输入和输出类型也可以通过反射自动提取出来。表值函数返回的表的类型取决于 TableFunction 类的泛型参数 T,不同于标量函数,表值函数的求值方法本身不包含返回类型,而是通过 collect(T) 方法来发送要输出的行。
  3. 在 Table API 中,表值函数是通过 .joinLateral(...) 或者 .leftOuterJoinLateral(...) 来使用的。joinLateral 算子会把外表(算子左侧的表)的每一行跟跟表值函数返回的所有行(位于算子右侧)进行 (cross)join。leftOuterJoinLateral 算子也是把外表(算子左侧的表)的每一行跟表值函数返回的所有行(位于算子右侧)进行(cross)join,并且如果表值函数返回 0 行也会保留外表的这一行。
  4. 在 SQL 里面用 JOIN 或者 以 ON TRUE 为条件的 LEFT JOIN 来配合 LATERAL TABLE(<TableFunction>) 的使用。
  5. 其实就是以前的UDTF函数
  1. package com.atguigu.flink.day10;
  2. import org.apache.flink.streaming.api.datastream.DataStreamSource;
  3. import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
  4. import org.apache.flink.table.annotation.DataTypeHint;
  5. import org.apache.flink.table.annotation.FunctionHint;
  6. import org.apache.flink.table.api.Table;
  7. import org.apache.flink.table.api.bridge.java.StreamTableEnvironment;
  8. import org.apache.flink.table.functions.ScalarFunction;
  9. import org.apache.flink.table.functions.TableFunction;
  10. import org.apache.flink.types.Row;
  11. import static org.apache.flink.table.api.Expressions.$;
  12. import static org.apache.flink.table.api.Expressions.call;
  13. /*
  14. hello hello hello 5
  15. hello 5
  16. hello world hello 5
  17. world 5
  18. atguigu hello hello atguigu 7
  19. ....
  20. */
  21. public class $03_FunctionTable {
  22. public static void main(String[] args) throws Exception {
  23. //获取流的执行环境
  24. StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
  25. env.setParallelism(1);
  26. //读取集合中的数据
  27. DataStreamSource<String> stream = env.fromElements(
  28. "hello hello","hello world","atguigu hello hello"
  29. );
  30. //获取表的执行环境
  31. StreamTableEnvironment tenv = StreamTableEnvironment.create(env);
  32. Table table = tenv.fromDataStream(stream);
  33. //1.在table api中使用
  34. //1.1内联的方式
  35. /*table
  36. .joinLateral(call(MySplit.class,$("f0")))
  37. .select($("f0"),$("word"),$("len"))
  38. .execute()
  39. .print();*/
  40. //1.2注册后使用
  41. /*tenv.createTemporaryFunction("my_split",MySplit.class);
  42. table
  43. .joinLateral(call("my_split",$("f0")))
  44. .select($("f0"),$("word"),$("len"))
  45. .execute()
  46. .print();*/
  47. //2.在SQL语句中使用
  48. //2.1先注册
  49. tenv.createTemporaryFunction("my_split",MySplit.class);
  50. //2.2再使用
  51. tenv.sqlQuery("select" + " f0, " +
  52. " word, " +
  53. " len " +
  54. "from " + table +
  55. " left join lateral table(my_split(f0)) on true"
  56. ).execute().print();
  57. /*tenv.sqlQuery("select" +
  58. " f0, " +
  59. " w, " +
  60. " l " +
  61. "from " + table +
  62. " left join lateral table(my_split(f0)) as T(w, l) on true"
  63. ).execute().print();*/
  64. }
  65. @FunctionHint(output = @DataTypeHint("row<word string, len int>"))
  66. public static class MySplit extends TableFunction<Row> {
  67. public void eval(String s){
  68. String[] words = s.split(" ");
  69. for (String word : words) {
  70. collect(Row.of(word,word.length()));
  71. }
  72. }
  73. }
  74. }

三.聚合函数

  1. 用户自定义聚合函数(User-Defined Aggregate Functions,UDAGGs)可以把一个表中的数据,聚合成一个标量值。用户定义的聚合函数,是通过继承AggregateFunction抽象类实现的。
  1. package com.atguigu.flink.day10;
  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 org.apache.flink.table.functions.AggregateFunction;
  8. import org.apache.flink.table.functions.ScalarFunction;
  9. import static org.apache.flink.table.api.Expressions.$;
  10. import static org.apache.flink.table.api.Expressions.call;
  11. /**
  12. * 变成大写字母的标量函数
  13. */
  14. public class $04_FunctionAgg {
  15. public static void main(String[] args) throws Exception {
  16. //获取流的执行环境
  17. StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
  18. env.setParallelism(1);
  19. //读取集合中的数据
  20. DataStreamSource<WaterSensor> stream = env.fromElements(
  21. new WaterSensor("sensor_1", 1000L, 10),
  22. new WaterSensor("sensor_1", 2000L, 20),
  23. new WaterSensor("sensor_2", 3000L, 30),
  24. new WaterSensor("sensor_1", 4000L, 40),
  25. new WaterSensor("sensor_1", 4000L, 50),
  26. new WaterSensor("sensor_2", 6000L, 60)
  27. );
  28. //获取表的执行环境
  29. StreamTableEnvironment tenv = StreamTableEnvironment.create(env);
  30. Table table = tenv.fromDataStream(stream);
  31. //1.在table api中使用
  32. //1.1内联的方式
  33. /*table
  34. .groupBy($("id"))
  35. .select($("id"),call(MyAvg.class,$("vc")).as("vc_avg"))
  36. .execute()
  37. .print();*/
  38. //1.2注册后使用
  39. /*tenv.createTemporaryFunction("my_avg",MyAvg.class);
  40. table
  41. .groupBy($("id"))
  42. .select($("id"),call("my_avg",$("vc")).as("vc_avg"))
  43. .execute()
  44. .print();*/
  45. //2.在SQL语句中使用
  46. //2.1先注册
  47. tenv.createTemporaryFunction("my_avg",MyAvg.class);
  48. //2.2再使用
  49. tenv.sqlQuery("select " +
  50. " id, " +
  51. " my_avg(vc) " +
  52. "from " + table +
  53. " group by id").execute().print();
  54. }
  55. public static class Avg{
  56. public Double sum = 0D;
  57. public Long count = 0L;
  58. public Double avg(){
  59. return sum / count;
  60. }
  61. }
  62. public static class MyAvg extends AggregateFunction<Double,Avg>{
  63. //返回最终的计算结果
  64. @Override
  65. public Double getValue(Avg acc) {
  66. return acc.avg();
  67. }
  68. //初始化累加器
  69. @Override
  70. public Avg createAccumulator() {
  71. return new Avg();
  72. }
  73. /**
  74. * 参数1:累加器 参数2:用户自定义的输入值
  75. * @param avg
  76. * @param vc
  77. */
  78. public void accumulate(Avg avg,Double vc){
  79. avg.sum += vc;
  80. avg.count++;
  81. }
  82. }
  83. }

四.表值聚合函数

  1. 自定义表值聚合函数(UDTAGG)可以把一个表(一行或者多行,每行有一列或者多列)聚合成另一张表,结果中可以有多行多列。
  1. package com.atguigu.flink.day10;
  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 org.apache.flink.table.functions.AggregateFunction;
  8. import org.apache.flink.table.functions.TableAggregateFunction;
  9. import org.apache.flink.util.Collector;
  10. import static org.apache.flink.table.api.Expressions.$;
  11. import static org.apache.flink.table.api.Expressions.call;
  12. /*
  13. 10 ..
  14. 第一 10
  15. 20
  16. 第一 20
  17. 第二 10
  18. 30
  19. 第一 30
  20. 第二 20
  21. */
  22. public class $05_FunctionTableAgg {
  23. public static void main(String[] args) throws Exception {
  24. //获取流的执行环境
  25. StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
  26. env.setParallelism(1);
  27. //读取集合中的数据
  28. DataStreamSource<WaterSensor> stream = env.fromElements(
  29. new WaterSensor("sensor_1", 1000L, 10),
  30. new WaterSensor("sensor_1", 2000L, 20),
  31. new WaterSensor("sensor_2", 3000L, 30),
  32. new WaterSensor("sensor_1", 4000L, 40),
  33. new WaterSensor("sensor_1", 4000L, 50),
  34. new WaterSensor("sensor_2", 6000L, 60)
  35. );
  36. //获取表的执行环境
  37. StreamTableEnvironment tenv = StreamTableEnvironment.create(env);
  38. Table table = tenv.fromDataStream(stream);
  39. //1.在table api中使用
  40. //1.1内联的方式
  41. table
  42. .groupBy($("id"))
  43. .flatAggregate(call(Top2Function.class,$("vc")))
  44. .select($("id"),$("level"),$("value"))
  45. .execute()
  46. .print();
  47. //1.2注册后使用
  48. /*tenv.createTemporaryFunction("top2",Top2Function.class);
  49. table
  50. .groupBy($("id"))
  51. .flatAggregate(call("top2",$("vc")))
  52. .select($("id"),$("level"),$("value"))
  53. .execute()
  54. .print();*/
  55. //2.在SQL语句中使用
  56. //不支持
  57. }
  58. public static class FirstSecond{
  59. public Integer first = 0;
  60. public Integer second = 0;
  61. }
  62. public static class Result{
  63. public String level;
  64. public Integer value;
  65. public Result(String level, Integer value) {
  66. this.level = level;
  67. this.value = value;
  68. }
  69. }
  70. public static class Top2Function extends TableAggregateFunction<Result,FirstSecond>{
  71. //初始化累加器
  72. @Override
  73. public FirstSecond createAccumulator() {
  74. return new FirstSecond();
  75. }
  76. //聚合
  77. public void accumulate(FirstSecond fs,Integer vc){
  78. if(vc > fs.first){
  79. fs.second = fs.first;
  80. fs.first = vc;
  81. }else if(vc > fs.second){
  82. fs.second = vc;
  83. }
  84. }
  85. //制表:通过out.collect()发射每行数据
  86. public void emitValue(FirstSecond fs, Collector<Result> out){
  87. out.collect(new Result("第一名",fs.first));
  88. if(fs.second>0){
  89. out.collect(new Result("第二名", fs.second));
  90. }
  91. }
  92. }
  93. }

第三章.SQL实现topN

1.介绍

  1. 目前仅 Blink 计划器支持 Top-N 。
  2. Flink 使用 OVER 窗口条件和过滤条件相结合以进行 Top-N 查询。利用 OVER 窗口的 PARTITION BY 子句的功能,Flink 还支持逐组 Top-N 。 例如,每个类别中实时销量最高的前五种产品。批处理表和流处理表都支持基于SQL的 Top-N 查询。
  3. 流处理模式需注意: TopN 查询的结果会带有更新。 Flink SQL 会根据排序键对输入的流进行排序;若 top N 的记录发生了变化,变化的部分会以撤销、更新记录的形式发送到下游。 推荐使用一个支持更新的存储作为 Top-N 查询的 sink 。另外,若 top N 记录需要存储到外部存储,则结果表需要拥有与 Top-N 查询相同的唯一键。

2.实现

需求:每隔30分钟统计最近1小时的热门商品top3,并把统计的结果写入到mysql中

思路:

  1. 按照商品id,窗口(hop)分组,计算点击量
  2. 使用over窗口:按照点击量进行排序,每个数据添加一个排名(row_number)
  3. 使用where 过滤出来topN where rn <= 3
  4. 把结果写入到mysql中 官方建议:把topN的结果写入到支持更新的数据库中
  1. 数据源

input/UserBehavior.csv

  1. 在Mysql中创建表
  1. CREATE DATABASE flink_sql;
  2. USE flink_sql;
  3. DROP TABLE IF EXISTS `hot_item`;
  4. CREATE TABLE `hot_item` (
  5. `w_end` timestamp NOT NULL,
  6. `item_id` bigint(20) NOT NULL,
  7. `item_count` bigint(20) NOT NULL,
  8. `rk` bigint(20) NOT NULL,
  9. PRIMARY KEY (`w_end`,`rk`)
  10. ) ENGINE=InnoDB DEFAULT CHARSET=utf8;
  1. 导入JDBC Connector依赖
  1. <dependency>
  2. <groupId>org.apache.flink</groupId>
  3. <artifactId>flink-connector-jdbc_${scala.binary.version}</artifactId>
  4. <version>${flink.version}</version>
  5. </dependency>
  1. 具体实现
  1. package com.atguigu.flink.day10;
  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. /**
  6. * 需求:每隔30分钟统计最近1小时的热门商品top3,并把统计的结果写入到mysql中
  7. */
  8. public class $08_TopN {
  9. public static void main(String[] args) {
  10. //获取流的执行环境
  11. StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
  12. env.setParallelism(1);
  13. //获取表的执行环境
  14. StreamTableEnvironment tenv = StreamTableEnvironment.create(env);
  15. //1.先建立一个动态表与数据源关联 事件时间
  16. tenv.executeSql(
  17. "create table ub(" +
  18. " user_id bigint, " +
  19. " item_id bigint, " +
  20. " category_id int, " +
  21. " behavior string, " +
  22. " ts bigint, " +
  23. " et as to_timestamp(from_unixtime(ts)), " +
  24. " watermark for et as et - interval '3' second " +
  25. ")with(" +
  26. " 'connector' = 'filesystem', " +
  27. " 'path' = 'input/UserBehavior.csv', " +
  28. " 'format' = 'csv' " +
  29. ")"
  30. );
  31. //tenv.sqlQuery("select * from ub").execute().print();
  32. //2.过滤pv数据,按照商品id 开窗,聚合
  33. Table t1 = tenv.sqlQuery(
  34. "select " +
  35. " item_id, " +
  36. " hop_start(et, interval '30' minute, interval '1' hour) stt, " +
  37. " hop_end(et, interval '30' minute, interval '1' hour) edt, " +
  38. " count(*) ct " +
  39. " from ub " +
  40. " where behavior='pv' " +
  41. " group by item_id, hop(et, interval '30' minute, interval '1' hour)"
  42. );
  43. tenv.createTemporaryView("t1",t1);
  44. //3.使用over窗口给每个聚合结果排序 row_number
  45. Table t2 = tenv.sqlQuery(
  46. "select" +
  47. " * , " +
  48. " row_number() over(partition by edt order by ct desc) rn " +
  49. "from t1 "
  50. );
  51. tenv.createTemporaryView("t2",t2);
  52. //4.过滤出topN
  53. Table t3 = tenv.sqlQuery("select" +
  54. " edt w_end, " +
  55. " item_id, " +
  56. " ct item_count, " +
  57. " rn rk " +
  58. " from t2 " +
  59. "where rn <= 3"
  60. );
  61. //t3.execute().print();
  62. // 5. 结果输出(sink:mysql)
  63. // 5.1 建立一张表与mysql关联
  64. tenv.executeSql("CREATE TABLE `hot_item` (\n" +
  65. " `w_end` timestamp ,\n" +
  66. " `item_id` bigint,\n" +
  67. " `item_count` bigint ,\n" +
  68. " `rk` bigint,\n" +
  69. " PRIMARY KEY (`w_end`,`rk`) not enforced\n" +
  70. ")with(" +
  71. " 'connector'='jdbc', " +
  72. " 'url'='jdbc:mysql://hadoop162:3306/flink_sql?useSSL=false', " +
  73. " 'table-name'='hot_item', " +
  74. " 'username'='root', " +
  75. " 'password'='aaaaaa' " +
  76. ") ");
  77. // 5.2 写入
  78. t3.executeInsert("hot_item");
  79. }
  80. }

day10[Flink SQL编程(下)] - 图1

第四章.双流join

在Flink中,支持两种方式的流的join:Window Join 和Interval Join

1.Window Join

窗口join会join具有相同的key并且处于同一个窗口的两个流的元素

  • 所有的窗口join都是inner join,意味着a 流中的元素如果在b 流中没有对应的,则a 流中这个元素就不会处理了(就是忽略掉了)
  • join成功后的元素会以所在窗口的最大时间作为时间戳,例如窗口[5,10),则元素会以9作为自己的时间戳
  1. package com.atguigu.flink.day10;
  2. import org.apache.flink.api.common.eventtime.WatermarkStrategy;
  3. import org.apache.flink.api.common.functions.JoinFunction;
  4. import org.apache.flink.api.java.tuple.Tuple3;
  5. import org.apache.flink.streaming.api.datastream.SingleOutputStreamOperator;
  6. import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
  7. import org.apache.flink.streaming.api.windowing.assigners.TumblingEventTimeWindows;
  8. import org.apache.flink.streaming.api.windowing.time.Time;
  9. /**
  10. * window join
  11. */
  12. public class $06_WindowJoin {
  13. public static void main(String[] args) throws Exception {
  14. StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
  15. env.setParallelism(1);
  16. SingleOutputStreamOperator<Tuple3<String, Long, Integer>> ds1 = env.fromElements(
  17. Tuple3.of("a", 1L, 1),
  18. Tuple3.of("a", 6L, 21),
  19. Tuple3.of("b", 3L, 12),
  20. Tuple3.of("a", 4L, 1111),
  21. Tuple3.of("b", 9L, 112)
  22. )
  23. .assignTimestampsAndWatermarks(
  24. WatermarkStrategy.<Tuple3<String, Long, Integer>>forMonotonousTimestamps()
  25. .withTimestampAssigner((data, ts) -> data.f1 * 1000L)
  26. );
  27. SingleOutputStreamOperator<Tuple3<String, Long, Integer>> ds2 = env.fromElements(
  28. Tuple3.of("a", 4L, 21),
  29. Tuple3.of("b", 7L, 111),
  30. Tuple3.of("d", 3L, 312),
  31. Tuple3.of("a", 9L, 9999),
  32. Tuple3.of("c", 3L, 110)
  33. )
  34. .assignTimestampsAndWatermarks(
  35. WatermarkStrategy.<Tuple3<String, Long, Integer>>forMonotonousTimestamps()
  36. .withTimestampAssigner((data, ts) -> data.f1 * 1000L)
  37. );
  38. ds1.join(ds2)
  39. .where(d1 -> d1.f0)
  40. .equalTo(d2 -> d2.f0)
  41. .window(TumblingEventTimeWindows.of(Time.seconds(5)))
  42. .apply(new JoinFunction<Tuple3<String, Long, Integer>, Tuple3<String, Long, Integer>, String>() {
  43. @Override
  44. public String join(Tuple3<String, Long, Integer> first, Tuple3<String, Long, Integer> second) throws Exception {
  45. //关联上的数据会进入这个方法
  46. //关联上:key一样,同一个窗口范围内
  47. return first + "<=========>" + second;
  48. }
  49. })
  50. .print();
  51. env.execute();
  52. }
  53. }

day10[Flink SQL编程(下)] - 图2

2.Interval Join

间隔流join(Interval Join)是指使用一个流的数据按照key去join另外一个流指定范围的数据

如下图:橙色的流去join绿色的流,范围是由橙色流的event-time + lower bound 和 event-time + upper bound来决定的

orangeElem.ts + lowerBound <= greenElem.ts <= orangeElem.ts + upperBound

day10[Flink SQL编程(下)] - 图3

  1. package com.atguigu.flink.day10;
  2. import org.apache.flink.api.common.eventtime.WatermarkStrategy;
  3. import org.apache.flink.api.common.functions.JoinFunction;
  4. import org.apache.flink.api.java.tuple.Tuple3;
  5. import org.apache.flink.streaming.api.datastream.KeyedStream;
  6. import org.apache.flink.streaming.api.datastream.SingleOutputStreamOperator;
  7. import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
  8. import org.apache.flink.streaming.api.functions.co.ProcessJoinFunction;
  9. import org.apache.flink.streaming.api.windowing.assigners.TumblingEventTimeWindows;
  10. import org.apache.flink.streaming.api.windowing.time.Time;
  11. import org.apache.flink.util.Collector;
  12. /**
  13. * Interval join
  14. * Intervaljoin实现的join效果,类似SQL里的innerjoin,取不到join不上的数据
  15. */
  16. public class $07_IntervalJoin {
  17. public static void main(String[] args) throws Exception {
  18. StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
  19. env.setParallelism(1);
  20. SingleOutputStreamOperator<Tuple3<String, Long, Integer>> ds1 = env.fromElements(
  21. Tuple3.of("a", 1L, 1),
  22. Tuple3.of("a", 6L, 21),
  23. Tuple3.of("b", 3L, 12),
  24. Tuple3.of("a", 4L, 1111),
  25. Tuple3.of("b", 9L, 112)
  26. )
  27. .assignTimestampsAndWatermarks(
  28. WatermarkStrategy.<Tuple3<String, Long, Integer>>forMonotonousTimestamps()
  29. .withTimestampAssigner((data, ts) -> data.f1 * 1000L)
  30. );
  31. SingleOutputStreamOperator<Tuple3<String, Long, Integer>> ds2 = env.fromElements(
  32. Tuple3.of("a", 4L, 21),
  33. Tuple3.of("b", 7L, 111),
  34. Tuple3.of("d", 3L, 312),
  35. Tuple3.of("a", 9L, 9999),
  36. Tuple3.of("c", 3L, 110)
  37. )
  38. .assignTimestampsAndWatermarks(
  39. WatermarkStrategy.<Tuple3<String, Long, Integer>>forMonotonousTimestamps()
  40. .withTimestampAssigner((data, ts) -> data.f1 * 1000L)
  41. );
  42. //1.先按照关联条件 keyBy
  43. KeyedStream<Tuple3<String, Long, Integer>, String> k1 = ds1.keyBy(d1 -> d1.f0);
  44. KeyedStream<Tuple3<String, Long, Integer>, String> k2= ds2.keyBy(d2 -> d2.f0);
  45. k1.intervalJoin(k2)
  46. .between(Time.seconds(-3),Time.seconds(2))
  47. .process(new ProcessJoinFunction<Tuple3<String, Long, Integer>, Tuple3<String, Long, Integer>, String>() {
  48. @Override
  49. public void processElement(Tuple3<String, Long, Integer> left, Tuple3<String, Long, Integer> right, Context ctx, Collector<String> out) throws Exception {
  50. out.collect(left + "<----------->" + right);
  51. }
  52. })
  53. .print();
  54. env.execute();
  55. }
  56. }

day10[Flink SQL编程(下)] - 图4

Interval Join原理:

  • 底层使用的 connect + keyby
  • 执行过程
  1. 两条流各初始化了一个状态,用来存储数据
  2. 先判断数据是否迟到,如果迟到,直接return,不处理
  3. 不管哪条流的数据来,都会存在自己的状态里
  4. 不管哪条流的数据来,都会遍历对方的状态,
    1. 如果在时间范围外,跳过
    2. 如果在时间范围内,join上—>发送给processElement方法 ——>我们实现的方法里拿到的就是join上的数据
  1. 清理数据的实现方式:注册一个定时器,到时间了就remove

第五章.海量数据实时去重

1.方案一:借用redis的Set

缺点:

  • 需要频繁连接redis
  • 如果数据量过大,对redis的内存也是一种压力

2.方案二:借用Flink的MapState

缺点:

  • 如果数据量过大,状态后端最好选择RocksDBStateBackend
  • 如果数据量过大, 对存储也有一定压力

3.方案三:使用布隆过滤器

布隆过滤器可以大大减少存储的数据的数据量

一.介绍

  1. 1.为什么需要布隆过滤器
  2. 如果想判断一个元素是不是在一个集合里,一般想到的是将集合中所有元素保存起来,然后通过比较确定。链表、树、散列表(又叫哈希表,Hash table)等等数据结构都是这种思路。
  3. 但是随着集合中元素的增加,我们需要的存储空间越来越大。同时检索速度也越来越慢,上述三种结构的检索时间复杂度分别为O(n),O(logn),O(1)。
  4. 布隆过滤器即可以解决存储空间的问题, 又可以解决时间复杂度的问题.
  5. 布隆过滤器的原理是,当一个元素被加入集合时,通过K个散列函数将这个元素映射成一个位数组中的K个点,把它们置为1。检索时,我们只要看看这些点是不是都是1就(大约)知道集合中有没有它了:如果这些点有任何一个0,则被检元素一定不在;如果都是1,则被检元素很可能在。这就是布隆过滤器的基本思想。
  1. 2.基本概念
  2. 布隆过滤器(Bloom Filter,下文简称BF)由Burton Howard Bloom在1970年提出,是一种空间效率高的概率型数据结构。它专门用来检测集合中是否存在特定的元素。
  3. 它实际上是一个很长的二进制向量和一系列随机映射函数。
  1. 3.实现原理
  2. 布隆过滤器的原理是,当一个元素被加入集合时,通过K个散列函数将这个元素映射成一个位数组中的K个点,把它们置为1。检索时,我们只要看看这些点是不是都是1就(大约)知道集合中有没有它了:如果这些点有任何一个0,则被检元素一定不在;如果都是1,则被检元素很可能在。这就是布隆过滤器的基本思想。
  3. BF是由一个长度为m比特的位数组(bit array)与k个哈希函数(hash function)组成的数据结构。位数组均初始化为0,所有哈希函数都可以分别把输入数据尽量均匀地散列。
  4. 当要插入一个元素时,将其数据分别输入k个哈希函数,产生k个哈希值。以哈希值作为位数组中的下标,将所有k个对应的比特置为1。
  5. 当要查询(即判断是否存在)一个元素时,同样将其数据输入哈希函数,然后检查对应的k个比特。如果有任意一个比特为0,表明该元素一定不在集合中。如果所有比特均为1,表明该集合有(较大的)可能性在集合中。为什么不是一定在集合中呢?因为一个比特被置为1有可能会受到其他元素的影响(hash碰撞),这就是所谓“假阳性”(false positive)。相对地,“假阴性”(false negative)在BF中是绝不会出现的。
  6. 下图示出一个m=18, k=3的BF示例。集合中的x、y、z三个元素通过3个不同的哈希函数散列到位数组中。当查询元素w时,因为有一个比特为0,因此w不在该集合中。

day10[Flink SQL编程(下)] - 图5

  1. 4.优点
  2. 一.不需要存储数据本身,只用比特表示,因此空间占用相对于传统方式有巨大的优势,并且能够保密数据;
  3. 二.时间效率也较高,插入和查询的时间复杂度均为O(K), 所以他的时间复杂度实际是O(1)
  4. 三.哈希函数之间相互独立,可以在硬件指令层面并行计算。
  1. 5.缺点
  2. 一.存在假阳性的概率,不适用于任何要求100%准确率的情境
  3. 二.只能插入和查询元素,不能删除元素,这与产生假阳性的原因是相同的。我们可以简单地想到通过计数(即将一个比特扩展为计数值)来记录元素数,但仍然无法保证删除的元素一定在集合中
  1. 6.使用场景
  2. 所以,BF在对查准度要求没有那么苛刻,而对时间、空间效率要求较高的场合非常合适.
  3. 另外,由于它不存在假阴性问题,所以用作“不存在”逻辑的处理时有奇效,比如可以用来作为缓存系统(如Redis)的缓冲,防止缓存穿

二.假阳性概率的计算

假阳性的概率.pdf

三.使用布隆过滤器实现去重

Flink已经内置了布隆过滤器的实现(使用的是google的Guava)

  1. package com.atguigu.flink.day10;
  2. import com.atguigu.flink.day07.pojo.UserBehavior;
  3. import org.apache.flink.shaded.guava18.com.google.common.hash.Funnels;
  4. import org.apache.flink.api.common.eventtime.WatermarkStrategy;
  5. import org.apache.flink.shaded.guava18.com.google.common.hash.BloomFilter;
  6. import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
  7. import org.apache.flink.streaming.api.functions.windowing.ProcessWindowFunction;
  8. import org.apache.flink.streaming.api.windowing.assigners.TumblingEventTimeWindows;
  9. import org.apache.flink.streaming.api.windowing.time.Time;
  10. import org.apache.flink.streaming.api.windowing.windows.TimeWindow;
  11. import org.apache.flink.util.Collector;
  12. import java.time.Duration;
  13. /**
  14. * 需求:指定时间范围内网站独立访客数(UV)的统计(使用布隆过滤器)
  15. */
  16. public class $09_BloomFilter {
  17. public static void main(String[] args) throws Exception {
  18. StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
  19. env.readTextFile("input/UserBehavior.csv")
  20. .map(line -> {
  21. String[] data = line.split(",");
  22. return new UserBehavior(
  23. Long.parseLong(data[0]),
  24. Long.parseLong(data[1]),
  25. Integer.parseInt(data[2]),
  26. data[3], Long.parseLong(data[4] ) * 1000
  27. );
  28. })
  29. .assignTimestampsAndWatermarks(
  30. WatermarkStrategy
  31. .<UserBehavior>forBoundedOutOfOrderness(Duration.ofSeconds(3))
  32. .withTimestampAssigner((ub,ts)-> ub.getTimestamp())
  33. )
  34. .filter(ub -> "pv".equals(ub.getBehavior()))
  35. .keyBy(ub -> ub.getBehavior())
  36. .window(TumblingEventTimeWindows.of(Time.minutes(30)))
  37. .process(new ProcessWindowFunction<UserBehavior, String, String, TimeWindow>() {
  38. @Override
  39. public void process(String s, Context context, Iterable<UserBehavior> elements, Collector<String> out) throws Exception {
  40. //1.创建一个布隆过滤器
  41. BloomFilter<Long> bf = BloomFilter.create(Funnels.longFunnel(), 1000000, 0.01);
  42. int uv = 0;
  43. for (UserBehavior element : elements) {
  44. if(bf.put(element.getUserId())){
  45. uv++;
  46. }
  47. }
  48. out.collect(context.window() +" " + uv);
  49. }
  50. })
  51. .print();
  52. env.execute();
  53. }
  54. }