Window Join 是把两条流分别划分到窗口中,然后只关联同一个窗口内、且满足 Join Key 的数据。

典型场景:

订单流 + 支付流

需求:

统计每个 5 分钟窗口内,订单和支付是否能够匹配。


1. 基本思想

假设使用 5 分钟滚动窗口:

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

订单和支付分别进入窗口:

order_id事件时间窗口
O00110:02[10:00, 10:05)
order_id支付时间窗口
O00110: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_startwindow_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_startwindow_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 版本;
  • 数据是否需要重新计算。

常见策略:

  1. 通过较大的 Watermark 延迟容忍乱序;
  2. 对严重迟到数据单独输出;
  3. 离线任务对历史窗口补算;
  4. 下游使用支持更新或撤回的存储;
  5. 对窗口结果设置版本或重新覆盖。

不要只看窗口大小,还要同时考虑:

窗口大小 + 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

更多推荐