TT Lab
开始
学习 学习路径 课程

Apache Flink — 用真正的引擎跑流处理

流式连接——记住什么、记多久

在 TT Lab 中继续学习

一句话总结

流式连接的区别在于记住哪一侧的行、记多久。没有时间条件的常规连接会永远持有两侧,区间连接在时间范围过去后就丢弃,而事件时间时态连接(temporal join)只持有右侧的版本,并把左侧的行与“该时刻的版本”中的一个连接起来。

为什么需要它

在批处理中,连接很简单。两张表都已经齐了,用一侧建哈希表,再扫描另一侧就完事了。在流中,两张表没有尽头。订单进来的瞬间,这笔订单的发货还不存在,而用户信息可能早就到了,也可能要很久之后才到。所以引擎必须把“可能与当前到来的行配对的对侧行”堆在某个地方,这堆东西就是状态。

问题在于什么时候丢弃。如果无限期地持有可能永远等不到配对的行,状态就会无限增长。可要是随便什么时候都丢,又会错过姗姗来迟的配对。Flink SQL 把连接分成多种,原因就在于此——引擎能安全丢弃什么,取决于查询对时间做出了什么承诺。

工作原理

三栏图。第一栏常规连接:订单和用户两侧的行都堆在状态里,无论哪一侧来了新行,都要与对侧的全部数据比对。第二栏区间连接:只有落在从订单时间起两小时宽的区带内的发货才会被连接,水位线越过区带末端后,该订单就会从状态中清除。第三栏时态连接:汇率像台阶一样变化的版本线,订单只与自己那个时刻有效的一级台阶连接。早于第一个版本的订单没有可连接的版本

常规连接(regular join) 最为自由。正如官方文档所述,一侧来了新行,就会与对侧的过去和未来的所有行比对。所以下午注册的用户在上午的订单也能连上——因为它不看时间。代价同样写在文档里:必须把两侧输入永远放在状态中。虽然可以用状态 TTL 来缩减,但那样结果可能出错。如果把执行计划导出成 JSON(COMPILE PLAN),会看到连接节点里写着 leftState 和 rightState 两个状态,TTL 都是 0 ms(不清除)。

区间连接(interval join) 需要一个等值条件,以及一个把两侧时间绑在一起的范围。s.ship_time BETWEEN o.order_time AND o.order_time + INTERVAL '2' HOUR 就是一个例子。输入必须是带时间属性的 append-only 表。时间属性基本单调递增,所以水位线一旦越过(订单时间 + 2 小时),就能确定不会再有与该订单配对的发货,于是从状态中清除。还有一个实测确认过的边界:BETWEEN 包含两端,延迟为 0 秒和恰好 7200 秒的发货可以连上,7201 秒的则不行。如果改成 LEFT JOIN,并把时间条件放在 ON 里,到窗口关闭时仍没有找到配对的订单,会与 null 一起输出一次。结果始终是 append-only。

事件时间时态连接 会把左侧(订单)的一行,与右侧版本表中“该时刻有效的版本”中的一个连接起来。语法是 SQL:2011 的 FOR SYSTEM_TIME AS OF o.order_time。要成为版本表,必须有主键和事件时间属性。像汇率文件这样 append-only 的 Source 不能声明主键,文档在这里给出了诀窍。按币种建立 ROW_NUMBER() OVER (PARTITION BY currency ORDER BY update_time DESC) = 1 的去重视图,优化器会推断出 currency 是主键,并把它作为版本视图使用。

CREATE TEMPORARY VIEW rates_v AS
SELECT currency, rate, update_time FROM (
  SELECT *, ROW_NUMBER() OVER (PARTITION BY currency ORDER BY update_time DESC) AS rn FROM rates)
WHERE rn = 1;

SELECT o.order_id, r.rate
FROM orders o JOIN rates_v FOR SYSTEM_TIME AS OF o.order_time AS r
  ON o.currency = r.currency;

实测结果是,订单与 update_time <= order_time 的版本中最晚的一个连接(如果汇率与订单在同一秒变化,则用新汇率)。早于第一个汇率的订单没有可连接的版本,所以在内连接中被排除了。正如文档所说,这种连接由两侧的水位线触发,右侧之后即使发生变化,也不会修正已经输出的结果。旧版本在不再需要时会从状态中清除。看计划 JSON 会发现,时态连接节点和区间连接节点里根本没有带 TTL 的 state 项——这两者不是靠 TTL,而是靠时间来清理状态。

如果同一个版本视图不加 FOR SYSTEM_TIME AS OF 就连接,那就只是普通的常规连接。每当汇率变化,过去的订单也会按新汇率重新计算,撤回(-U)和更新(+U)就会大量涌出,最终结果是按最后一个汇率换算的值。

连接 记住什么 是否修正结果 清除状态的依据
常规 两侧全部 会修正(输入有更新时) 只有 TTL(可能丢失正确性)
区间 时间范围内的行 不修正 水位线越过范围末端
事件时间时态 右侧需要的版本 不修正 水位线越过后不再需要的版本

在现场相遇的样子

最常见的事故是“只是给订单附加商品信息,状态却涨了好几个月”。原因几乎总是常规连接。商品表再小,订单一侧也会永远堆积。如果维度信息只需要那个时点的值,时态连接才是合适的工具。

第二种是销售额重算事故。把汇率或价格表用常规连接挂上去,价格变化的瞬间,过去订单的金额就会悄悄改变。收到“仪表板上昨天的销售额今天变了”的报告时,先看连接类型。在本实验中,把同样的 USD 订单用两种方式换算,合计确实会不同。

第三种是区间连接的边界。“两小时内发货”是写成 < 还是 BETWEEN,决定了恰好两小时发出的货算不算。运营指标的定义书和 SQL 中的不等号,至少要核对一次。

下一项实验要做什么

定义四个 Source,用常规连接汇总各等级的订单。用区间连接连上两小时内的发货,再用外连接形式的区间连接找出没有按时发货的订单。把汇率做成版本视图,用时态连接连上订单时刻的汇率,再看同一个视图用常规连接后合计有何不同。把三种连接的执行计划导出成 JSON,比较其中的状态项,最后把数字整理成报告。