Flink sql有几种join方式
Flink SQL 的 Join 可以从两个角度理解:
- SQL 语义:Inner / Left / Right / Full / Cross Join
- 流计算实现与时间语义:Regular Join、Interval Join、Window Join、Temporal Join、Lookup Join
实际开发中第二种分类更重要,因为它决定了状态大小、数据延迟、是否支持维表、能否处理历史版本。
一、按普通 SQL 语义分类
1. Inner Join(内连接)
只保留两边能匹配上的记录。
SELECT *
FROM orders o
JOIN users u
ON o.user_id = u.user_id;
2. Left Outer Join(左外连接)
保留左表所有数据;右表匹配不到时,右侧字段为 NULL。
SELECT *
FROM orders o
LEFT JOIN users u
ON o.user_id = u.user_id;
流计算中用于“订单流关联维表/用户流”很常见。
3. Right Outer Join(右外连接)
保留右表全部记录。
SELECT *
FROM orders o
RIGHT JOIN users u
ON o.user_id = u.user_id;
虽然 SQL 层面可能支持,但实际流处理里通常不常写。因为一般会交换表顺序,改写成 Left Join,使语义更直观。
4. Full Outer Join(全外连接)
两边数据都保留,匹配不到的一侧补 NULL。
SELECT *
FROM orders o
FULL OUTER JOIN users u
ON o.user_id = u.user_id;
流式场景使用时要非常谨慎:Flink 不知道未来是否还会有匹配记录,因此可能需要长时间甚至无限期维护状态。
5. Cross Join(笛卡尔积)
每条左表数据和每条右表数据组合。
SELECT *
FROM a
CROSS JOIN b;
等价于没有有效关联条件的 Join。
流式作业中基本应避免,因为数据量可能迅速膨胀:
左流 1 万条 × 右流 1 万条 = 1 亿条组合
有时小型静态维表广播关联会看起来类似 Cross Join,但应使用 Lookup Join 或 Broadcast Hint 等更合适的方式,而不要真的做大规模笛卡尔积。
二、Flink 流式开发中常见的 Join 类型
1. Regular Join(常规流流 Join)
也叫普通 Join。两个动态表按关联键连接:
SELECT
o.order_id,
o.user_id,
u.user_name
FROM orders o
JOIN users u
ON o.user_id = u.user_id;
例如:
订单流 JOIN 用户变更流
特点
- 不限制时间范围;
- 两侧到达的数据都可能与未来数据匹配;
- 理论上需要保存两边历史状态;
- 如果没有 TTL 或窗口限制,状态可能持续增长。
例如:
A 流:user_id = 1,10:00 到达
B 流:user_id = 1,18:00 到达
Flink 仍可能认为两者需要匹配,因此要保留状态。
适用场景
- 两个 CDC 表进行实时关联;
- 双流都需要持续更新;
- 数据量可控,或已配置合理的状态 TTL。
但对于大规模无限流,通常更建议使用带时间范围的 Join。
2. Interval Join(区间 Join / 时间范围 Join)
两个流按 Key 关联,并且要求事件时间处在一个明确范围内。
SELECT
o.order_id,
p.pay_id
FROM orders o
JOIN payments p
ON o.order_id = p.order_id
AND p.rowtime BETWEEN o.rowtime - INTERVAL '10' MINUTE
AND o.rowtime + INTERVAL '5' MINUTE;
含义:
支付时间在订单时间前 10 分钟到后 5 分钟之间,才认为匹配。
也可以写成:
AND o.rowtime BETWEEN p.rowtime - INTERVAL '5' MINUTE
AND p.rowtime + INTERVAL '10' MINUTE
特点
- 需要事件时间属性,如
rowtime; - 有明确时间边界;
- Watermark 推进后,过期状态可以清理;
- 比普通 Regular Join 更适合无限流。
典型场景
- 下单后 30 分钟内是否支付;
- 点击后 10 分钟内是否下单;
- 两个系统的同一业务事件在一定时间范围内对账;
- 请求日志和响应日志关联。
3. Window Join(窗口 Join)
先对两张表按相同窗口切分,再在同一个窗口内 Join。
例如 5 分钟滚动窗口:
SELECT
o.window_start,
o.window_end,
o.order_id,
p.pay_id
FROM TABLE(
TUMBLE(TABLE orders, DESCRIPTOR(rowtime), INTERVAL '5' MINUTES)
) o
JOIN TABLE(
TUMBLE(TABLE payments, DESCRIPTOR(rowtime), INTERVAL '5' MINUTES)
) p
ON o.order_id = p.order_id
AND o.window_start = p.window_start
AND o.window_end = p.window_end;
特点
- 两张表通常要使用相同类型、相同大小、相同偏移量的窗口;
- 窗口结束、水位线推进后,Flink 可以释放对应窗口状态;
- 状态边界更明确;
- 常用于窗口聚合后或窗口内事件匹配。
常见窗口类型
TUMBLE:滚动窗口,例如每 5 分钟;HOP:滑动窗口,例如每 5 分钟计算一次、覆盖过去 1 小时;CUMULATE:累积窗口;SESSION:会话窗口,使用时需更谨慎。
与 Interval Join 的区别
- Interval Join:以某条事件为中心定义前后时间范围;
- Window Join:同一个固定/滑动窗口内的数据互相匹配。
例如“订单后 30 分钟内支付”更适合 Interval Join;“同一个 5 分钟窗口内的订单和支付统计”更偏向 Window Join。
4. Temporal Join(时态 Join / 版本表 Join)
Temporal Join 用于关联具有历史版本的动态维表。
典型场景:
订单流 + 商品价格 CDC 流
商品价格会不断修改。你希望订单关联的不是“现在最新价格”,而是:
订单发生时,该商品有效的价格。
示例:
SELECT
o.order_id,
o.product_id,
o.order_time,
p.price
FROM orders o
JOIN product_price p
FOR SYSTEM_TIME AS OF o.order_time
ON o.product_id = p.product_id;
其中:
orders:事实流;product_price:带主键和变更历史的版本表;o.order_time:订单的事件时间;FOR SYSTEM_TIME AS OF:表示按订单发生时刻查询维度版本。
要求和特点
通常要求维表:
- 有主键(Primary Key);
- 是支持变更的动态表,例如 CDC 来源;
- 有版本时间语义;
- 上游最好能提供正确的 update / delete 变更。
典型场景
- 订单关联下单时商品价格;
- 交易关联交易发生时客户等级;
- 业务事件关联当时的组织架构;
- 广告点击关联点击时广告配置。
这和数仓中的 SCD Type 2 / 拉链表思想非常接近。
5. Lookup Join(查找 Join / 维表 Join)
Lookup Join 通常是事实流去查询一个外部维表。
SELECT
o.order_id,
o.user_id,
d.user_name,
d.user_level
FROM orders o
LEFT JOIN dim_user FOR SYSTEM_TIME AS OF o.proc_time AS d
ON o.user_id = d.user_id;
其中 dim_user 可以来自:
- MySQL;
- Redis;
- HBase;
- JDBC;
- Elasticsearch;
- Paimon、Iceberg、Hudi 等支持 Lookup 的表;
- 自定义 Connector。
特点
- 左侧通常是实时事实流;
- 右侧通常是外部维度表;
- 通过 Key 查找右表数据;
- 常支持缓存、异步 I/O;
- 多数场景使用处理时间
PROCTIME()。
示例:定义处理时间
CREATE TABLE orders (
order_id STRING,
user_id STRING,
amount DECIMAL(10, 2),
ts TIMESTAMP(3),
proc_time AS PROCTIME()
) WITH (...);
然后:
SELECT
o.order_id,
o.amount,
u.user_name
FROM orders AS o
LEFT JOIN dim_user FOR SYSTEM_TIME AS OF o.proc_time AS u
ON o.user_id = u.user_id;
注意
这种写法通常表示:
订单到达 Flink 并执行关联时,去拿“当时查到的最新维度值”。
它不一定等于“订单业务发生时的维度版本”。
因此:
- Lookup Join 常用于实时补充当前维度信息;
- Event-time Temporal Join 更适合还原业务事件发生时的历史维度属性。
三、如何选择
| 需求 | 更适合的 Join |
|---|---|
| 两张无界流按 Key 持续关联 | Regular Join,但需控制状态 |
| 两个事件必须在限定时间内匹配 | Interval Join |
| 同一窗口内的数据关联 | Window Join |
| 实时订单查询 MySQL / Redis 当前维度 | Lookup Join |
| 订单关联发生时有效的价格、等级、组织 | Event-time Temporal Join |
| 普通静态或批表关联 | 普通 Inner / Outer Join |
| 无条件组合两张表 | Cross Join,通常避免 |
四、一个实务上的判断口诀
可以粗略记成:
流 + 流,且有时间关系 -> Interval Join / Window Join
流 + 外部当前维表 -> Lookup Join
流 + CDC 版本维表 -> Temporal Join
两个流只是按 Key 硬关联 -> Regular Join,注意状态膨胀
五、特别注意:普通 Join 在流上可能导致状态无限增长
下面这种 SQL 很容易出问题:
SELECT *
FROM orders o
JOIN payments p
ON o.order_id = p.order_id;
如果 orders 和 payments 都是无界流,且没有时间范围,Flink 理论上需要持续保存已经到达的数据,因为未来任何时刻都可能出现匹配记录。
生产中至少应考虑:
- 使用 Interval Join 或 Window Join;
- 配置状态 TTL;
- 判断业务是否允许超时后不再匹配;
- 将一侧改为 Lookup 维表;
- 使用 CDC / Temporal Join 表达版本关联。
所以,“Flink SQL 有几种 Join”若按工程实践回答,最常需要掌握的是:
Regular Join、Interval Join、Window Join、Temporal Join、Lookup Join。
它们内部再结合INNER、LEFT OUTER等 SQL Join 语义使用。
更多推荐
所有评论(0)