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 | // 获取StreamTableEnvironment. |
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 | # create a Table from a Table API query |
1.6 视图 view
1.7 window aggregate与group aggregate区别
参考自Apache Flink 零基础入门(九):Flink SQL 编程实践
1.7.1 Group Aggregate的例子
这是一个group aggregate, 内存中累计每个分类的数据,有变化时,update sink
1 | SELECT psgCnt, COUNT(*) AS cnt |
1.7.2 Window Aggregate
1 | SELECT |
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 | SELECT DATE_FORMAT(ts, 'yyyy-MM-dd'), shop_id, COUNT(*) as pv |
当然,如果 TTL 配置地太小,可能会清除掉一些有用的状态和数据,从而导致数据精确性地问题。这也是用户需要权衡地一个参数。
[参考文献]