Flink sql固定时间窗口join
Window Join 是把两条流分别划分到窗口中,然后只关联同一个窗口内、且满足 Join Key 的数据。
典型场景:
订单流 + 支付流
需求:
统计每个 5 分钟窗口内,订单和支付是否能够匹配。
1. 基本思想
假设使用 5 分钟滚动窗口:
[10:00, 10:05)
[10:05, 10:10)
[10:10, 10:15)
订单和支付分别进入窗口:
| order_id | 事件时间 | 窗口 |
|---|---|---|
| O001 | 10:02 | [10:00, 10:05) |
| order_id | 支付时间 | 窗口 |
|---|---|---|
| O001 | 10:04 | [10:00, 10:05) |
因为:
order_id相同;window_start相同;window_end相同;
所以两条记录可以 Join。
如果支付发生在 10:07,它属于 [10:05, 10:10),就不能和订单 O001 匹配。
2. 基本语法
Flink SQL 中通常通过 Windowing TVF 定义窗口:
SELECT
o.window_start,
o.window_end,
o.order_id,
o.amount,
p.pay_amount,
p.pay_time
FROM TABLE(
TUMBLE(
TABLE orders,
DESCRIPTOR(order_time),
INTERVAL '5' MINUTE
)
) AS o
JOIN TABLE(
TUMBLE(
TABLE payments,
DESCRIPTOR(pay_time),
INTERVAL '5' MINUTE
)
) AS p
ON o.window_start = p.window_start
AND o.window_end = p.window_end
AND o.order_id = p.order_id;
这里有三个关键 Join 条件:
o.window_start = p.window_start
AND o.window_end = p.window_end
AND o.order_id = p.order_id
其中:
order_id:业务关联键;window_start:窗口开始时间;window_end:窗口结束时间。
窗口字段不是普通业务字段,但必须用于保证两边数据来自同一个窗口。
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'
);
5 分钟窗口 Join
SELECT
o.window_start,
o.window_end,
o.order_id,
o.user_id,
o.amount,
p.pay_id,
p.pay_amount
FROM TABLE(
TUMBLE(
TABLE orders,
DESCRIPTOR(order_time),
INTERVAL '5' MINUTE
)
) AS o
JOIN TABLE(
TUMBLE(
TABLE payments,
DESCRIPTOR(pay_time),
INTERVAL '5' MINUTE
)
) AS p
ON o.window_start = p.window_start
AND o.window_end = p.window_end
AND o.order_id = p.order_id;
4. 为什么需要 Watermark?
Window Join 通常使用事件时间,因此需要 Watermark 判断窗口是否已经基本完成。
例如:
WATERMARK FOR order_time AS order_time - INTERVAL '5' SECOND
表示允许最多约 5 秒的乱序数据。
对于窗口:
[10:00, 10:05)
只有当 Watermark 推进到接近或超过 10:05 时,Flink 才能判断:
- 这个窗口是否完成;
- 是否还可能有迟到数据;
- 什么时候可以清理窗口状态。
没有 Watermark,Flink 很难可靠地判断窗口何时结束,状态也可能无法及时清理。
5. window_start 和 window_end 为什么必须 Join?
假设只写:
ON o.order_id = p.order_id
那么订单 O001 和支付 O001 在不同窗口中也可能被匹配:
订单:10:02,属于 [10:00, 10:05)
支付:10:08,属于 [10:05, 10:10)
这可能不符合“同一窗口内匹配”的业务要求。
因此应写成:
ON o.order_id = p.order_id
AND o.window_start = p.window_start
AND o.window_end = p.window_end
这也是 Window Join 和普通 Regular Join 的重要区别。
6. Tumbling、Hopping 和 Cumulate Window Join
6.1 Tumbling Window Join
滚动窗口,窗口之间不重叠:
TUMBLE(
TABLE orders,
DESCRIPTOR(order_time),
INTERVAL '5' MINUTE
)
窗口:
[10:00, 10:05)
[10:05, 10:10)
[10:10, 10:15)
适合:
- 每 5 分钟订单和支付匹配;
- 每小时事件对账;
- 固定时间段数据关联。
6.2 Hopping Window Join
滑动窗口,窗口之间可能重叠:
HOP(
TABLE orders,
DESCRIPTOR(order_time),
INTERVAL '1' MINUTE,
INTERVAL '10' MINUTE
)
含义:
- 每 1 分钟产生一个窗口;
- 每个窗口覆盖 10 分钟。
窗口可能是:
[10:00, 10:10)
[10:01, 10:11)
[10:02, 10:12)
一条数据可能进入多个窗口,因此 Join 结果也可能出现多份。
适合:
- 最近 10 分钟的数据匹配;
- 滑动监控;
- 高灵敏度的实时统计。
注意:两边必须使用一致的窗口参数,否则 window_start 和 window_end 很难正确匹配。
6.3 Cumulate Window Join
累积窗口会在窗口范围内逐步扩大:
CUMULATE(
TABLE orders,
DESCRIPTOR(order_time),
INTERVAL '5' MINUTE,
INTERVAL '1' HOUR
)
例如:
[10:00, 10:05)
[10:00, 10:10)
[10:00, 10:15)
...
[10:00, 11:00)
这种窗口更适合累积指标。由于同一数据可能进入多个累积窗口,Join 结果也需要注意重复和语义。
7. Window Join 和 Interval Join 的区别
两者都能限制时间范围,但语义不同。
Window Join
要求两条数据落在同一个固定窗口:
订单 10:02
支付 10:04
同属 [10:00, 10:05),可以匹配
如果支付 10:06:
订单 10:02 -> [10:00, 10:05)
支付 10:06 -> [10:05, 10:10)
即使只相差 4 分钟,也不能匹配。
Interval Join
直接定义两条事件之间允许的时间差:
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
订单 10:02,支付 10:06,间隔 4 分钟,可以匹配。
因此:
| 需求 | 选择 |
|---|---|
| 必须在同一固定窗口内匹配 | Window Join |
| 只要求事件时间差在范围内 | Interval Join |
| 订单后 30 分钟内支付 | Interval Join |
| 同一个 5 分钟批次内对账 | Window Join |
8. Window Join 和 Regular Join 的区别
Regular Join
SELECT *
FROM orders o
JOIN payments p
ON o.order_id = p.order_id;
没有时间范围限制,未来任何时间到达的数据理论上都可能匹配。状态可能持续增长。
Window Join
ON o.order_id = p.order_id
AND o.window_start = p.window_start
AND o.window_end = p.window_end
数据被限制在窗口中,窗口结束后可以清理状态,状态边界更明确。
可以简单理解为:
Regular Join:只按 Key 匹配
Window Join:按 Key + 窗口匹配
9. 窗口 Join 中的 Left Join
如果希望保留没有支付的订单,可以使用 Left Join:
SELECT
o.window_start,
o.window_end,
o.order_id,
o.amount,
p.pay_id,
p.pay_amount
FROM TABLE(
TUMBLE(
TABLE orders,
DESCRIPTOR(order_time),
INTERVAL '5' MINUTE
)
) AS o
LEFT JOIN TABLE(
TUMBLE(
TABLE payments,
DESCRIPTOR(pay_time),
INTERVAL '5' MINUTE
)
) AS p
ON o.window_start = p.window_start
AND o.window_end = p.window_end
AND o.order_id = p.order_id;
匹配不到支付时:
pay_id = NULL
pay_amount = NULL
但生产使用时需要注意版本差异和结果更新语义。流式 Left Join 可能先输出一条右侧为空的结果,之后支付数据到达后再产生更新或撤回结果,具体表现取决于 Flink 版本、Connector 和下游是否支持更新流。
如果业务要求:
窗口结束后最终判断哪些订单没有支付
通常需要在窗口关闭后再做结果处理,而不是简单依赖实时 Left Join 的中间结果。
10. 迟到数据怎么处理?
假设窗口是:
[10:00, 10:05)
Watermark 已经超过 10:05 后,订单 10:02 才到达,这就是迟到数据。
处理方式取决于:
- Watermark 延迟;
- 是否配置允许迟到;
- 下游是否支持更新;
- 使用的窗口 TVF 和 Flink 版本;
- 数据是否需要重新计算。
常见策略:
- 通过较大的 Watermark 延迟容忍乱序;
- 对严重迟到数据单独输出;
- 离线任务对历史窗口补算;
- 下游使用支持更新或撤回的存储;
- 对窗口结果设置版本或重新覆盖。
不要只看窗口大小,还要同时考虑:
窗口大小 + Watermark 延迟 + 业务允许的迟到时间
11. 常见错误
错误一:两边窗口配置不一致
例如左边:
TUMBLE(..., INTERVAL '5' MINUTE)
右边:
TUMBLE(..., INTERVAL '10' MINUTE)
两边的窗口边界不一致,通常无法得到正确匹配。
错误二:忘记关联窗口字段
错误:
ON o.order_id = p.order_id
正确:
ON o.order_id = p.order_id
AND o.window_start = p.window_start
AND o.window_end = p.window_end
错误三:使用处理时间却期待事件时间结果
如果使用处理时间,结果取决于数据何时到达 Flink,而不是业务事件实际发生时间。
如果订单和支付存在乱序或延迟,通常应优先考虑事件时间和 Watermark。
错误四:认为 Window Join 只会输出最终一条结果
窗口内数据不断到达时,Join 结果可能持续产生或更新。下游需要根据结果类型支持:
- Append;
- Upsert;
- Update Before / Update After;
- Retract。
不能默认所有结果都是只追加。
12. 一句话总结
Window Join 的核心 SQL 结构是:
FROM WindowTVF(left_table) l
JOIN WindowTVF(right_table) r
ON l.window_start = r.window_start
AND l.window_end = r.window_end
AND l.business_key = r.business_key
它适合:
两条流中的数据必须落在同一个时间窗口内,且业务键也匹配的场景。
如果你的需求是“订单发生后一定时间内支付”,不要机械地使用 Window Join。由于固定窗口边界可能把时间上很接近的事件分到不同窗口,这种情况通常更适合 Interval Join。
更多推荐
所有评论(0)