Flink SQL 的 Join 可以从两个角度理解:

  1. SQL 语义:Inner / Left / Right / Full / Cross Join
  2. 流计算实现与时间语义: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;

如果 orderspayments 都是无界流,且没有时间范围,Flink 理论上需要持续保存已经到达的数据,因为未来任何时刻都可能出现匹配记录。

生产中至少应考虑:

  • 使用 Interval Join 或 Window Join;
  • 配置状态 TTL;
  • 判断业务是否允许超时后不再匹配;
  • 将一侧改为 Lookup 维表;
  • 使用 CDC / Temporal Join 表达版本关联。

所以,“Flink SQL 有几种 Join”若按工程实践回答,最常需要掌握的是:

Regular Join、Interval Join、Window Join、Temporal Join、Lookup Join
它们内部再结合 INNERLEFT OUTER 等 SQL Join 语义使用。

更多推荐