「我的数据空间」实时计算实践笔记 · Flink SQL 系列
引言
无论在OLAP领域还是OLTP领域,多表Join都是业务所必备的。在OLTP场景中,日常事务的处理需要使用到Join操作;OLAP场景由于数据量大、字段多,数据通常被分为事实表和维度表以星型或雪花模型组成,那么数据查询中多表Join的操作更是不能少。对于离线计算而言,经过数据库领域多年的积累,Join 语义以及实现已经十分成熟,然而对于近年来刚兴起的实时领域 Streaming SQL 来说 Join 却处于刚起步的状态。
Join的本质是分别从N(N >= 1)张表中获取不同字段进行拼接,从而获得完整的数据。目前SQL中Join的种类可以分为
- CROSS Join,交叉连接,计算表的笛卡尔积。
- INNER Join,内连接,计算表的交集。
- OUTER Join
- LEFT,返回左表所有行,右表不存在补NULL。
- RIGHT,返回右表的所有行,左表不存在补NULL。
- FULL,返回左右表的并集,不存在的补NULL。
- SELF Join,自连接。
在Batch SQL的模式中,表中的数据是一个有限集,Join的实现会依赖数据集的缓存。然而在Streaming SQL的模式中,数据集是无限的。那么如何在无限的数据流中实现表的Join操作,使得特定时刻的输出结果与Batch模式结果相同?想要解决这个问题,首先需要先了解动态表与连续查询的概念。
动态表与连续查询
在Flink的官方文档中,专门有一页介绍了动态表与连续查询的概念,在此不赘述,只是简单介绍下相关概念,以及为什么无限流数据可以进行查询和Join。
这里需要提到一个概念——流表对偶(duality)性。对偶性是描述导致相同的物理结果,表面上不同的理论之间的对应关系。那么流表对偶性即描述数据流和数据表虽然属于不同的理论,但他们承载的SQL却能够通过联系拥有相同的物理结果。那么如何关联流和表呢,下面的例子就能够很好地展示。
例:通过MySQL主从复制机制,了解流表对偶性。
了解MySQL的同学应该都知道,Binlog是MySQL实现主从复制的关键,Binlog会记录CREATE、ALTER TABLE等和INSERT、UPDATE等操作。如果选择模式为row-based,那么对表的操作会以数据行为单位进行记录。记录在BInlog中的数据会发给从服务器,并在从服务器上进行表的还原。
根据上图,理解动态表会很容易。动态表,就是会随时间所变化的表,随着数据地不断流入,表中字段的内容也会发生变化。相对于静态表而言,针对动态表的查询略有区别。针对动态表的查询并不会终止,而是一直在进行,类似于数据库中存在的物化视图。在进行连续查询时,Flink首先会将changelog stream(INSERT、UPDATE、DELETE操作)转化为一张动态表,并在此动态表上进行持续查询。再将查询的结果再生成一张结果动态表,并依据此结果动态表将数据一条一条的地写入下游。
双流Join
了解了上文的两个概念后,理解Streaming模式下的Join就会变得很简单。Flink在双流Join这一类中支持两种Join方式,分别是Regular Join和Interval Join。其原理基本相同,Regular Join可以认为就是Batch模式Join的另一种表现形式,在同一时刻其结果与Batch模式完全相同。Interval Join为了解决Regular Join的数据无限增长的问题,引入了时间窗口概念。
Regular Join
该模式是最基本的Join模式,对数据的更新或改变对Join两边都是可见的,且随时间变化能够影响全局结果。有的同学可能会有疑问,流式Join数据是一条一条被处理的,该条数据影响之后的Join结果很容易理解,那么如何影响之前的Join结果呢?看完下面的内容,你可能就会有答案。
上图清晰地展示了Join的处理流程。来自两边的数据(L-Event、R-Event)进入Join算子后先会被分别更新到L-State和R-State,接着L-Event会和R-State中的结果进行Join操作,输出Join结果发到下游;同样,R-Event则会去和L-State进行Join,并把相关结果发往下游。
假设我们有两张source表,分别是Orders表和Shipments表,数据源都是changelog。
Orders:+--------+----------+id|price|+--------+----------+01|10|+--------+----------+02|20|+--------+----------+03|30|+--------+----------+---------------------------------------------------Shipments+--------+----------+id|type|+--------+----------+01|AA|+--------+----------+03|AB|+--------+----------+05|BB|+--------+----------+通过id来进行Join,DML如下:selecto.id,o.price,s.typefromOrders oLEFTJoinShipments sono.id=s.id;我们假设两张表中数据的流入顺序为:
Orders | Shipments ------------------------------ 1. + (01, 10) | ------------------------------ 2. + (02, 20) | ------------------------------ 3. | + (01, AA) ------------------------------ 4. | + (03, AB) ------------------------------ 5. + (03, 30) | ------------------------------ 6. | + (05, BB) 注:+号代表此数据为Insert那么Join后流出的数据为:
isINSERT|id|price|type-------------------------------------------1.true|01|10|NULL-------------------------------------------2.true|02|20|NULL-------------------------------------------3.false|01|10|NULLtrue|01|10|AA-------------------------------------------4.-------------------------------------------5.true|03|30|AB-------------------------------------------6.上面的情况较为简单,都是append模式,那么如果存在retract的数据,情况又会变得复杂很多。假设在上述结果之后,我们有第7条数据进入:
Orders | Shipments ------------------------------ 7. | - (01, AA) 注:- 代表delete那么对应第7条输出为:
isINSERT|id|price|type-------------------------------------------7.false|01|10|AAtrue|01|10|NULL从上面的逻辑不难看出,因为有着retract的存在,使得Flink能够及时纠正Join结果,使结果与batch模式保持一致。双流Join的具体逻辑在flink-table-runtime-blink 模块下的StreamingJoinOperator类中,感兴趣的同学可以深入研究。目前Flink支持INNER/LEFT/RIGHT/FULL Join。
SEMI Join AND ANTI Join
在正常的Join外,还存在着一类较为特殊的Join方式,分别是SEMI Join和ANTI Join,这两种Join的特殊之处在于他们只返回左表的列数据,并不将右表的数据做输出,只是用右表数据对左表进行过滤。下图直观地展现了 SEMI Join 和 ANTI Join 输出结果的异同。
-- SEMI Join:SELECT*FROMEmployeeWHEREDeptNameIN(SELECTDeptNamefromDept);-- ANTI Join:SELECT*FROMEmployeeWHEREDeptNameNOTIN(SELECTDeptNamefromDept);虽然Regular Join能够像Batch模式那样满足我们的需求,但是他也有着致命缺点,即两张表的数据都需要缓存在state中,在unbound数据的情况下,占用的资源会无限增长。为了解决这样的问题,且保证Join的核心逻辑,Flink引入了Interval Join。
Interval Join
为了解决Regular Join数据持续增长的问题,Flink在Interval Join中引入了时间窗口的概念,窗口外的数据会被Flink清理,这极大缓解了资源的占用。Interval Join的时间语义既可以是Event Time,也可以是Processing Time,Flink会根据选择的时间语义来维护窗口。
Orders:+----------+----------+----------+|order_time|id|price|+----------+----------+----------+|XXXXXXXXXX|01|10|+----------+----------+----------+|XXXXXXXXXX|02|20|+----------+----------+----------+|XXXXXXXXXX|03|30|+----------+----------+----------+---------------------------------------------------Shipments+----------+----------+----------+|ship_time|id|type|+----------+----------+----------+|XXXXXXXXX|01|AA|+----------+----------+----------+|XXXXXXXXX|03|AB|+----------+----------+----------+|XXXXXXXXX|05|BB|+----------+----------+----------+例如上面的两张表,如果我们需要在4小时窗口内查看商品的发货信息,可以用下文DML实现:
CREATETABLEOrders(order_timeBIGINT,idINT,priceDOUBLE,WATERMARKFORToTIMESTAMP(order_time)asevent_time-INTERVAL'60'SECOND)with(..........)CREATETABLEShipments(ship_timeBIGINT,idINT,typeVARCHAR,WATERMARKFORToTIMESTAMP(ship_time)asevent_time-INTERVAL'60'SECOND)with(..........)INSERTINTOprintSELECT*FROMOrders o,Shipments sWHERE0.id=s.idANDs.ship_timeBETWEENo.order_timeANDo.order_time+INTERVAL'4'HOUR;从上图可以看出,根据Shipments表的watermark,order_time小于ship_time - 4 hour的数据可以被丢弃;根据Orders表的watermark,ship_time小于order_time的数据会被丢弃,Flink会依据watermark来进行过期数据的清理,将空间维持在合适的范围。
维表Join
虽然Interval Join能够解决资源问题,但是也给表绑定了时间界,超出时间界限的数据需要被丢弃。为了支持数据量不大,变化不频繁表的一类Join场景,Flink引入了两种维表Join模式,分别是Temporal Table Join和Temporal Table Function Join。这两种Join模式的主要应用场景是为了补全事实表(probe table)数据的额外字段,通常维度表(build table)里的字段很少发生变化。从数据量上来看,维度表一般数据量较小。通常维表Join的逻辑是基于Hash Join来实现的,Flink中的维表Join逻辑也不例外。Hash Join的Join逻辑正常会分为两步,第一步,将维度表按Key散列,建立哈希表;第二步,用事实表中的RowKey去哈希表中探测,再输出结果,所以通常又会将事实表叫做probe table,维度表叫做build table。
Temporal Table Join
按维表加载方式,实现逻辑主要分为两类。第一类是全量加载,例如HDFS,Flink会将HDFS文件全量加载进内存,进行Join操作时再去内存的cache匹配;第二类为部分加载,例如JDBC、Hbase、Redis等,会根据probe table当前Row数据的Key去数据库查询,再依据情况辅之以缓存逻辑。在使用该模式进行Join操作时,需要先在probe table中指定Processing time列,具体的Join写法如下:
insertintoSinkselecto.amount,o.currency,l.rate,o.amount*l.ratefromOrderojoinLatestRateFORSYSTEM_TIMEASOFo.proctimeaslONo.currency=l.currency;需要注意的是,目前Temporal Table Join仅支持INNER JOIN和LEFT JOIN。
Temporal Table Join的缺点是不能指定Event time作为时间语义,只支持Processing time,那就是说,无论probe table处在哪个时间段,都会和build table中最新的数据进行Join。
Temporal Table Function Join
Temporal Table Function模式主要是通过UDTF来实现probe流和Temporal table的Join。需要注意的是,这里的left input(probe table)需要是append-only table,right input(build table)需要有主键和用于版本化的字段(通常是时间字段)。在具体实现逻辑中,Flink会将左表数据和右表数据按Key分别保存到leftState和rightState中,格式都是MapState<Long, BaseRow>。左右state不同的是,leftState的主键是一个递增序列,rightState则以时间列作为主键。在Join逻辑中,首先会遍历左边的状态state,并提取元素中的时间列,用时间列去排好序的rightState中进行binary search,查找rightState.rowTime <= leftState.rowTime的第一条数据,匹配上就发往下游,同时在leftState状态中清除。从上面的匹配逻辑可以看出,目前Temporal Table Function Join仅支持INNER Join。
Temporal Table Function Join模式在SQL语言方面仅支持probe table的DDL和Join逻辑的DML,暂时还不支持以SQL来创建Temporal Table Function。想要使用该模式首先需要借助Table API创建Temporal Table Function,再写DML进行Join。
总结
实时领域的Streaming Join因为不能像Batch Join那样缓存完整数据集,所以需要给缓存设定基于时间的清理机制和限定Join涉及的数据范围。Flink SQL针对这些差异,从双流Join和维表Join两个方向设计,推出了多种Join模式来满足日常业务需求。
本文收录于「我的数据空间」技术库——一套可私有化部署的数据平台(数据集成 / 实时计算 / 数据湖 / 湖仓查询 / 智能问数)。产品介绍见[我的数据空间官网]。