flink Interval Join 入门
Interval Join 用于关联两条流中,满足以下条件的数据:
- 业务关联键相同;
- 两条数据的事件时间差在指定区间内。
典型场景:
订单流 + 支付流
业务规则:
订单发生后 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 Join | Interval 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_id | order_time |
|---|---|
| O001 | 10:00 |
支付流:
| pay_id | order_id | pay_time |
|---|---|---|
| P001 | O001 | 10:05 |
| P002 | O001 | 10:10 |
如果两条支付都在订单后的 30 分钟内,那么结果可能是:
| order_id | pay_id |
|---|---|
| O001 | P001 |
| O001 | P002 |
因此如果业务要求“每个订单只保留第一笔支付”,需要额外做去重或 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 分钟内发生的支付,可以和订单匹配。
更多推荐
所有评论(0)