0%

FlinkSQL与动态表

1 Flink SQL

1.1 Fink SQL核心功能一览

  • SELECT FROM WHERE

  • GROUP BY /HAVING

  • Regular AGG

  • Time-windowed AGG(TUMBLE,HOP,SESSION)

  • Regular JOIN(INNER,LEFT/RIGHT/FULL)

  • Time-windowed JOIN(INNER,LEFT/RIGHT/FULL)

  • Temporal Table Join

  • 大量内置函数(150+)

  • LIKE, EXTRACT, TIMESTAMPADD, MD5, AVG..

  • 丰富的类型系统:POJO,Map,Array,Row,NestedType

  • 自定义函数(scalar,table,aggregate)

  • CEP on SQL(复杂事件处理)

https://flink.apache.org/2020/07/28/flink-sql-demo-building-e2e-streaming-application.html

http://wuchong.me/blog/2019/08/20/flink-sql-training/

https://github.com/ververica/flink-sql-cookbook

1.2 Retract mode 和 Append mode

toAppendStream 只支持insert
toRetractStream 其余模式都可以

如果动态表仅只有Insert操作,即之前输出的结果不会被更新,则使用该模式。如果更新或删除操作使用追加模式会失败报错,始终可以使用此模式。返回值是boolean类型。它用true或false来标记数据的插入和撤回,返回true代表数据插入,false代表数据的撤回。

使用flinkSQL处理实时数据当我们把表转化成流的时候,需要用toAppendStream与toRetractStream这两个方法。稍不注意可能直接选择了toAppendStream。

始终可以使用此模式。返回值是boolean类型。它用true或false来标记数据的插入和撤回,返回true代表数据插入,false代表数据的撤回。

当我们使用的sql语句包含:count() group by时,必须使用缩进模式

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
// 获取StreamTableEnvironment.
StreamTableEnvironment tableEnv = ...;
// 包含两个字段的表(String name, Integer age)
Table table = ...
// 将表转为DataStream,使用Append Mode追加模式,数据类型为Row
DataStream<Row> dsRow = tableEnv.toAppendStream(table, Row.class);
// 将表转为DataStream,使用Append Mode追加模式,数据类型为定义好的TypeInformation
TupleTypeInfo<Tuple2<String, Integer>> tupleType = new TupleTypeInfo<>(
Types.STRING(),
Types.INT());
DataStream<Tuple2<String, Integer>> dsTuple =
tableEnv.toAppendStream(table, tupleType);
// 将表转为DataStream,使用的模式为Retract Mode撤回模式,类型为Row
// 对于转换后的DataStream<Tuple2<Boolean, X>>,X表示流的数据类型,
// boolean值表示数据改变的类型,其中INSERT返回true,DELETE返回的是false
DataStream<Tuple2<Boolean, Row>> retractStream
tableEnv.toRetractStream(table, Row.class);

1.2.1 如何实现回退更新?

flink-connector-jdbc 最终使用的是SQL引擎的upsert语法:

1
insert into tbl() values ( ),( ) on duplicate key update

可以看下 MySQL/upsert 一节

对应flink 源码

1
org.apache.flink.connector.jdbc.dialect.MySQLDialect#getUpsertStatement

1.3 keyedBy与group by区别

1.4 SQL解析工具

hive使用了antlr3实现了自己的HQL,
Flink使用Apache Calcite,
而Calcite的解析器是使用JavaCC实现的,
Spark2.x以后采用了antlr4实现自己的解析器,
Presto也是使用antlr4。

1.5 通过Table api创建表

1
2
# create a Table from a Table API query
tapi_result = table_env.from_path("table1").select(...)

1.6 视图 view

1.7 window aggregate与group aggregate区别

参考自Apache Flink 零基础入门(九):Flink SQL 编程实践

1.7.1 Group Aggregate的例子

这是一个group aggregate, 内存中累计每个分类的数据,有变化时,update sink

1
2
3
4
SELECT psgCnt, COUNT(*) AS cnt 
FROM Rides
WHERE isInNYC(lon, lat)
GROUP BY psgCnt;

1.7.2 Window Aggregate

1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
SELECT 
toAreaId(lon, lat) AS area, -- toAreaId为UDF
TUMBLE_END(rideTime, INTERVAL '5' MINUTE) AS window_end,
--滚动窗口的结束时间:
-- ① 可以时间不同吗?
-- ② 窗口可以别名吗?
COUNT(*) AS cnt
FROM Rides
WHERE isInNYC(lon, lat) and isStart
-- isStart 为boolean
GROUP BY
toAreaId(lon, lat),
TUMBLE(rideTime, INTERVAL '5' MINUTE)
-- 定义窗口
HAVING COUNT(*) >= 5;
-- 5分钟超过5次才输出

1.7.3 Window Aggregate 与 Group Aggregate 的区别

Window Aggregate 是当window结束时才输出,其输出的结果是最终值,不会再进行修改,其输出流是一个 Append 流。而 Group Aggregate 是每处理一条数据,就输出最新的结果,其结果是在不断更新的,就好像数据库中的数据一样,其输出流是一个 Update 流。

window 由于有 watermark ,可以精确知道哪些窗口已经过期了,所以可以及时清理过期状态,保证状态维持在稳定的大小。而 Group Aggregate 因为不知道哪些数据是过期的,所以状态会无限增长,这对于生产作业来说不是很稳定,所以建议对 Group Aggregate 的作业配上 State TTL 的配置。

例如统计每个店铺每天的实时PV,那么就可以将 TTL 配置成 24+ 小时,因为一天前的状态一般来说就用不到了。

1
2
3
SELECT  DATE_FORMAT(ts, 'yyyy-MM-dd'), shop_id, COUNT(*) as pv
FROM T
GROUP BY DATE_FORMAT(ts, 'yyyy-MM-dd'), shop_id

当然,如果 TTL 配置地太小,可能会清除掉一些有用的状态和数据,从而导致数据精确性地问题。这也是用户需要权衡地一个参数。

[参考文献]


  1. 在数据流中使用SQL查询:Apache Flink中的动态表的持续查询

  2. Flink Table API & SQL编程指南(1)

  3. Flink SQL金竹学习文档