news 2026/9/7 11:28:55

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

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
FlinkSQL 常用 Join 方式:Regular / Interval / 维表 Join 怎么选

「我的数据空间」实时计算实践笔记 · 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模式来满足日常业务需求。


本文收录于「我的数据空间」技术库——一套可私有化部署的数据平台(数据集成 / 实时计算 / 数据湖 / 湖仓查询 / 智能问数)。产品介绍见[我的数据空间官网]。

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/9/7 11:28:12

AI公司盈利背后:大模型推理成本优化与工程化实践

一家从成立之初就把“深度学习”写进基因的公司&#xff0c;能够在 2026 年上半年首次实现盈利&#xff0c;这件事放在整个 AI 行业里&#xff0c;都是一个非常值得拆解的信号。大家通常看到的是“盈利”这个财务结果&#xff0c;但作为长期关注大模型工程落地的开发者&#xf…

作者头像 李华
网站建设 2026/9/5 16:38:33

GNSS 高级篇 04 信号体制:4.3 信号强化技术解析

GNSS 高级篇 4 信号体制 4.3 信号强化技术解析 前面两节讲了信号"是什么"和"怎么共存"&#xff0c;这一节我们讲信号"怎么变强"——在真实世界里&#xff0c;GNSS 信号面临着干扰、多径、欺骗等各种威胁&#xff0c;信号体制和接收机技术是怎么应…

作者头像 李华
网站建设 2026/9/5 9:58:35

从Anthropic IPO传闻看Claude API的接入与网络排查

这几天 AI 圈最热的讨论&#xff0c;除了模型能力本身&#xff0c;还有一个消息值得开发者关注&#xff1a;Anthropic 被传正在考虑 IPO&#xff0c;并且可能允许内部人分批套现。这个信息如果只看新闻标题&#xff0c;感觉是财经频道的事&#xff0c;但对你我这种每天调 Claud…

作者头像 李华
网站建设 2026/9/4 15:11:12

VMware Workstation Pro 从零到自动化:虚拟机安装、配置、快照与网络详解

这次我们直接看虚拟机。很多开发者和学生都需要在 Windows 主机上再跑一个 Linux 或 Windows 测试环境&#xff0c;VMware Workstation Pro 就是目前用得最多的桌面级虚拟机软件之一。它解决的问题很具体&#xff1a;不需要重装系统、不需要双系统切换&#xff0c;就能在同一台…

作者头像 李华
网站建设 2026/9/3 3:18:45

微信小程序云开发实战:从0到1搭建校园生活圈

简介&#xff1a;这是一套面向高校学生与小程序开发初学者的实战型校园生活服务类应用源码&#xff0c;基于微信小程序云开发技术构建&#xff0c;无需自建服务器即可快速部署表白墙、失物招领、兼职信息发布及闲置物品买卖四大核心功能&#xff0c;切实解决校园内信息互通与轻…

作者头像 李华
网站建设 2026/9/5 21:02:35

基于51单片机的电子万年历设计:从硬件选型到代码调试

简介&#xff1a;本资源是一份面向高校电子类、计算机类专业本科生的单片机硬件课程设计实践材料&#xff0c;聚焦51系列单片机&#xff08;如AT89C51&#xff09;开发电子万年历系统&#xff0c;完整覆盖时间显示、闰年自动校准、按键调节与断电记忆等核心功能&#xff0c;适用…

作者头像 李华