Interval Join 用于关联两条流中,满足以下条件的数据:

  1. 业务关联键相同;
  2. 两条数据的事件时间差在指定区间内。

典型场景:

订单流 + 支付流

业务规则:

订单发生后 30 分钟内完成支付,认为订单支付成功。

这类需求更适合 Interval Join,而不是固定窗口 Join。


1. 基本语义

假设订单时间为:

order_time = 10:00

要求支付时间在订单前 10 分钟到订单后 30 分钟之间:

pay_time BETWEEN order_time - 10 MINUTE
                  AND order_time + 30 MINUTE

那么:

支付时间是否匹配
09:55
10:00
10:20
10:30
10:31

Interval Join 的核心是:

业务 Key 相等 + 时间位于指定区间

2. Flink SQL 基本写法

SELECT
    o.order_id,
    o.user_id,
    o.order_time,
    o.amount,
    p.pay_id,
    p.pay_time,
    p.pay_amount
FROM orders AS o
JOIN payments AS p
ON o.order_id = p.order_id
AND p.pay_time BETWEEN
        o.order_time - INTERVAL '10' MINUTE
    AND o.order_time + INTERVAL '30' MINUTE;

含义是:

o.order_id = p.order_id

并且:

o.order_time - 10分钟
<= p.pay_time
<= o.order_time + 30分钟

3. 完整建表示例

订单流

CREATE TABLE orders (
    order_id   STRING,
    user_id    STRING,
    amount     DECIMAL(10, 2),
    order_time TIMESTAMP(3),

    WATERMARK FOR order_time AS
        order_time - INTERVAL '5' SECOND
) WITH (
    'connector' = 'kafka',
    'topic' = 'orders',
    'format' = 'json'
);

支付流

CREATE TABLE payments (
    pay_id     STRING,
    order_id   STRING,
    pay_amount DECIMAL(10, 2),
    pay_time   TIMESTAMP(3),

    WATERMARK FOR pay_time AS
        pay_time - INTERVAL '5' SECOND
) WITH (
    'connector' = 'kafka',
    'topic' = 'payments',
    'format' = 'json'
);

Interval Join

SELECT
    o.order_id,
    o.user_id,
    o.order_time,
    o.amount,
    p.pay_id,
    p.pay_time,
    p.pay_amount
FROM orders AS o
JOIN payments AS p
ON o.order_id = p.order_id
AND p.pay_time BETWEEN
        o.order_time - INTERVAL '10' MINUTE
    AND o.order_time + INTERVAL '30' MINUTE;

这里的两个时间字段都应该是事件时间属性,并且通常需要定义 Watermark。


4. 为什么需要 Watermark?

Interval Join 需要知道:

某个时间范围内的数据是不是已经基本到齐了?

例如:

订单时间:10:00
允许支付范围:[09:50, 10:30]

如果当前只收到一条 10:10 的支付记录,Flink 不能立刻认为没有其他匹配,因为未来可能还会到达 10:20、10:29 的支付记录。

当两边的 Watermark 都推进到足够晚的时间后,Flink 才能清理已经不可能再产生匹配的数据。

因此,Interval Join 通常依赖:

WATERMARK FOR event_time AS event_time - INTERVAL '5' SECOND

Watermark 中的 5 秒表示允许一定程度的乱序。


5. BETWEEN 的边界含义

下面的写法:

p.pay_time BETWEEN
        o.order_time - INTERVAL '10' MINUTE
    AND o.order_time + INTERVAL '30' MINUTE

通常表示包含两端:

下界 <= pay_time <= 上界

也可以写成明确的比较条件:

AND p.pay_time >= o.order_time - INTERVAL '10' MINUTE
AND p.pay_time <= o.order_time + INTERVAL '30' MINUTE

实践中建议统一团队规范,避免边界事件在不同任务中出现不同解释。


6. 两种常见写法

订单后固定时间内支付

SELECT
    o.order_id,
    p.pay_id
FROM orders o
JOIN payments p
ON o.order_id = p.order_id
AND p.pay_time >= o.order_time
AND p.pay_time <= o.order_time + INTERVAL '30' MINUTE;

表示:

order_time <= pay_time <= order_time + 30分钟

这是最常见的“事件发生后关联”场景。


两个事件时间差不超过 5 分钟

SELECT
    a.device_id,
    a.event_time AS a_time,
    b.event_time AS b_time
FROM event_a a
JOIN event_b b
ON a.device_id = b.device_id
AND b.event_time BETWEEN
        a.event_time - INTERVAL '5' MINUTE
    AND a.event_time + INTERVAL '5' MINUTE;

表示两条流中同一设备的事件时间差不超过 5 分钟。


7. Interval Join 和 Window Join 的区别

这是最容易混淆的地方。

Window Join:同一个固定窗口

例如 5 分钟滚动窗口:

[10:00, 10:05)
[10:05, 10:10)

订单发生在 10:04,支付发生在 10:06:

订单 -> [10:00, 10:05)
支付 -> [10:05, 10:10)

即使两者只相差 2 分钟,也不能匹配。


Interval Join:直接比较事件时间差

订单 10:04,支付 10:06:

pay_time - order_time = 2分钟

如果规则是订单后 30 分钟内支付,那么可以匹配。

对比:

维度Window JoinInterval Join
时间规则必须在同一个窗口时间差满足区间
是否受固定窗口边界影响
适合订单后 30 分钟内支付一般更适合
适合固定 5 分钟批次对账更适合也可以但不够直观
是否需要窗口字段 Join不需要
是否需要 Watermark通常需要通常需要

记忆方式:

同一批次 -> Window Join
时间差范围 -> Interval Join

8. Interval Join 和 Regular Join 的区别

Regular Join

SELECT *
FROM orders o
JOIN payments p
ON o.order_id = p.order_id;

只按业务键关联,没有时间限制。

对于两条无界流,理论上未来任何时刻的数据都可能匹配,因此需要长期保存状态,状态可能持续增长。


Interval Join

ON o.order_id = p.order_id
AND p.pay_time BETWEEN
        o.order_time
    AND o.order_time + INTERVAL '30' MINUTE

数据只在指定时间范围内匹配,超出范围后就可以清理状态。

因此:

Interval Join 本质上是给 Regular Join 增加了明确的时间边界。


9. 一条数据可能匹配多条数据

Interval Join 不保证一对一。

例如:

order_idorder_time
O00110:00

支付流:

pay_idorder_idpay_time
P001O00110:05
P002O00110:10

如果两条支付都在订单后的 30 分钟内,那么结果可能是:

order_idpay_id
O001P001
O001P002

因此如果业务要求“每个订单只保留第一笔支付”,需要额外做去重或 Top 1 处理,例如:

  • 按支付时间排序;
  • 使用 ROW_NUMBER()
  • 根据支付状态过滤;
  • 在下游做幂等更新。

不要把 Interval Join 自动理解成一对一 Join。


10. Left Join 能不能用?

某些 Flink 版本和执行计划支持 Interval Left Join,可以保留左侧没有匹配的数据:

SELECT
    o.order_id,
    o.order_time,
    p.pay_id,
    p.pay_time
FROM orders o
LEFT JOIN payments p
ON o.order_id = p.order_id
AND p.pay_time BETWEEN
        o.order_time
    AND o.order_time + INTERVAL '30' MINUTE;

但这里有一个实时语义问题:

订单刚到时,支付可能还没有到达。Flink 如果立即输出:

order_id = O001, pay_id = NULL

之后支付到达,可能还要输出更新结果或撤回结果。

所以要确认下游是否支持:

  • UPDATE_BEFORE
  • UPDATE_AFTER
  • Retract;
  • Upsert;
  • 按主键更新。

如果只是要判断“30 分钟后仍未支付”,通常需要引入定时器、窗口结束逻辑或延迟结果设计,而不是简单依赖实时 Left Join。


11. 时间字段必须注意类型

Interval Join 通常要求关联时间字段是时间属性,例如:

TIMESTAMP(3)

或带时区的时间类型,具体取决于 Flink 版本。

推荐直接定义:

event_time TIMESTAMP(3),
WATERMARK FOR event_time AS event_time - INTERVAL '5' SECOND

不要在 Join 条件中对时间字段做复杂函数转换,例如:

DATE_FORMAT(order_time, ...)
CAST(order_time AS STRING)

这可能导致无法识别为有效的时间区间条件,或者无法使用 Interval Join 的优化执行方式。


12. 状态和性能

Interval Join 的状态大小主要取决于:

数据流速率 × 时间区间大小 × Key 分布

例如:

p.pay_time BETWEEN
    o.order_time - INTERVAL '10' MINUTE
AND o.order_time + INTERVAL '30' MINUTE

Flink 可能需要保留覆盖这段时间范围的数据。

如果每秒有 10 万条数据,关联范围是 40 分钟,状态规模就可能很大。

优化方向:

缩小时间范围

业务允许时,不要无必要地设置很大的区间:

INTERVAL '30' MINUTE

通常比:

INTERVAL '7' DAY

更容易控制状态。

先过滤数据

SELECT *
FROM payments
WHERE pay_status = 'SUCCESS'

减少进入 Join 的数据量。

控制 Key 分布

如果大量数据集中在同一个 Key,可能发生数据倾斜。例如:

order_id = UNKNOWN
user_id = 0
device_id = EMPTY

这些异常 Key 会导致大量匹配。

合理设置 Watermark

Watermark 太保守,状态释放慢;太激进,迟到数据容易被丢弃或产生错误结果。


13. 典型应用场景

订单和支付匹配

同一个 order_id
支付时间在订单后 30 分钟内

点击和下单归因

同一个 user_id
下单发生在点击后 24 小时内

请求和响应匹配

同一个 request_id
响应时间在请求之后 5 分钟内

日志和告警关联

同一个 device_id
异常日志和告警时间相差不超过 1 分钟

多系统事件对账

同一个业务流水号
两边事件时间相差不超过 10 分钟

14. 选择 Interval Join 时的检查清单

使用前确认:

  • 两边是否都是流表或动态表?
  • 是否有明确的业务关联键?
  • 两边是否都有事件时间字段?
  • 是否定义了 Watermark?
  • 时间关系是否可以表达成固定区间?
  • 是否允许一对多匹配?
  • 是否需要处理迟到数据?
  • 下游是否支持更新或撤回结果?
  • 时间范围是否足够小,可以控制状态?
  • 是否存在严重数据倾斜 Key?

总结

Interval Join 的核心语法是:

SELECT ...
FROM left_stream l
JOIN right_stream r
ON l.business_key = r.business_key
AND r.event_time BETWEEN
        l.event_time - INTERVAL '前置时间'
    AND l.event_time + INTERVAL '后置时间';

它解决的是:

两条流中,业务键相同且事件时间差在指定范围内的数据关联问题。

最典型的例子:

ON o.order_id = p.order_id
AND p.pay_time BETWEEN
        o.order_time
    AND o.order_time + INTERVAL '30' MINUTE

表示:

同一个订单,在下单后 30 分钟内发生的支付,可以和订单匹配。

更多推荐