Skip to content

FlinkSQL 常用 Join 方式:Regular / Interval / 维表 Join 怎么选

引言

无论在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中的数据会发给从服务器,并在从服务器上进行表的还原。

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结果呢?看完下面的内容,你可能就会有答案。

Regular Join:两侧各维护 state,互相探测

上图清晰地展示了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。

sql
Orders:
+--------+----------+
   id    |   price  |
+--------+----------+
   01    |    10    |
+--------+----------+
   02    |    20    |
+--------+----------+
   03    |    30    |
+--------+----------+

---------------------------------------------------
Shipments
+--------+----------+
   id    |   type   |
+--------+----------+
   01    |    AA    |
+--------+----------+
   03    |    AB    |
+--------+----------+
   05    |    BB    |
+--------+----------+

通过id来进行Join,DML如下:
select o.id, o.price, s.type 
from Orders o LEFT Join Shipments s
on o.id = s.id;

我们假设两张表中数据的流入顺序为:

markdown
     Orders    |    Shipments
------------------------------
1. + (01, 10)  |
------------------------------
2. + (02, 20)  |
------------------------------
3.             |   + (01, AA)
------------------------------
4.             |   + (03, AB)
------------------------------
5. + (03, 30)  |
------------------------------
6.             |   + (05, BB)
注:+号代表此数据为Insert

那么Join后流出的数据为:

ruby
   isINSERT   |   id   |   price   |   type   
-------------------------------------------
1.    true    |   01   |   10      |   NULL
-------------------------------------------
2.    true    |   02   |   20      |   NULL
-------------------------------------------
3.    false   |   01   |   10      |   NULL
      true    |   01   |   10      |   AA
-------------------------------------------
4.     
-------------------------------------------
5.    true    |   03   |   30      |   AB
-------------------------------------------
6.

上面的情况较为简单,都是append模式,那么如果存在retract的数据,情况又会变得复杂很多。假设在上述结果之后,我们有第7条数据进入:

markdown
     Orders   |    Shipments
------------------------------
7.            |  - (01, AA)
注:- 代表delete

那么对应第7条输出为:

ruby
   isINSERT   |   id   |   price   |   type   
-------------------------------------------
7.    false   |   01   |   10      |   AA
      true    |   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 与 ANTI Join 的输出差异

sql
-- SEMI Join:
SELECT * FROM Employee WHERE DeptName IN (SELECT DeptName from Dept);
sql
-- ANTI Join:
SELECT * FROM Employee WHERE DeptName NOT IN (SELECT DeptName from Dept);

虽然Regular Join能够像Batch模式那样满足我们的需求,但是他也有着致命缺点,即两张表的数据都需要缓存在state中,在unbound数据的情况下,占用的资源会无限增长。为了解决这样的问题,且保证Join的核心逻辑,Flink引入了Interval Join。

Interval Join

为了解决Regular Join数据持续增长的问题,Flink在Interval Join中引入了时间窗口的概念,窗口外的数据会被Flink清理,这极大缓解了资源的占用。Interval Join的时间语义既可以是Event Time,也可以是Processing Time,Flink会根据选择的时间语义来维护窗口。

ruby
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实现:

sql
CREATE TABLE Orders(
    order_time BIGINT,
    id INT,
    price DOUBLE,
    WATERMARK FOR ToTIMESTAMP(order_time) as event_time - INTERVAL '60' SECOND 
)with(
    ..........
)
CREATE TABLE Shipments(
    ship_time BIGINT,
    id INT,
    type VARCHAR,
    WATERMARK FOR ToTIMESTAMP(ship_time) as event_time - INTERVAL '60' SECOND
)with(
    ..........
)
INSERT INTO print
SELECT * FROM 
Orders o, Shipments s 
WHERE 0.id = s.id AND
s.ship_time BETWEEN o.order_time AND o.order_time + INTERVAL '4' HOUR;

Interval Join 的有效区间与过期清理

从上图可以看出,根据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写法如下:

sql
insert into Sink
select 
o.amount, o.currency, l.rate, o.amount * l.rate 
from Order o
join LatestRate FOR SYSTEM_TIME AS OF o.proctime as l 
ON o.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模式来满足日常业务需求。


本文是「我的数据空间」实时计算实践笔记的一篇。平台把这类链路做成了托管作业(Flink SQL 流作业、实时集成、Iceberg 表管理与维护),支持私有化部署与 OEM 合作 —— 见核心能力总览技术博客,交流合作 QQ:1559851993。