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

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

用三种连接关联订单

在 TT Lab 中继续学习

目标

对同一个订单流,分别用常规连接、区间连接和事件时间时态连接连上用户、发货和汇率,并通过输出和执行计划确认各种连接记住了什么、是否会修正结果。

为什么重要

流式连接不知道配对何时到来,所以要把对侧的行堆在状态里。没有时间条件就永远堆,给出时间范围就堆到水位线越过为止,如果是版本表,则只持有需要的版本。选错连接类型,状态会无限增长,或者过去的结果会被悄悄改变。本实验的评分器不会查询集群,而是把你保存的 sql-client 输出和计划 JSON,与直接从源 CSV 计算出的值进行核对。

步骤

  1. 用 flink-up 启动集群,在 /root/flink/joins/ddl.sql 中定义四个 Source(users · orders · shipments · rates)。在 orders、shipments、rates 的时间列上设置水位线。把以批处理方式统计四张表行数的 /root/flink/joins/count.sql 的输出保存到 /root/flink/joins/count.out(列名 tbl、n)。
  2. 在 /root/flink/joins/regular.sql 中写一个流式查询:把 orders 和 users 按 user_id 连接,输出各等级的 orders(订单数)和 amount(合计),并把输出保存到 /root/flink/joins/regular.out。
  3. 在 /root/flink/joins/interval.sql 中写区间连接,连上从订单时间起两小时内(含两端)发出的货,把 order_id、ship_id、delay_s(秒)保存到 /root/flink/joins/interval.out。
  4. 在 /root/flink/joins/unshipped.sql 中写外连接形式的区间连接,输出两小时内没有发货的订单的 order_id、order_time,并把输出保存到 /root/flink/joins/unshipped.out。
  5. 用 rates 创建版本视图 rates_v,用时态连接连上订单时刻的汇率,把输出 order_id、currency、rate、amount_krw(amount × rate)的 /root/flink/joins/temporal.sql 的输出保存到 /root/flink/joins/temporal.out。
  6. 把同一个 rates_v 不加 FOR SYSTEM_TIME AS OF 连接,把输出各币种 orders 和 total_krw 的 /root/flink/joins/latest.sql 的输出保存到 /root/flink/joins/latest.out。
  7. 用 COMPILE PLAN 把三种连接(第 2、3、5 步)的执行计划导出为 /root/flink/joins/regular-plan.json、/root/flink/joins/interval-plan.json、/root/flink/joins/temporal-plan.json。
  8. 在 /root/flink/joins/report.json 中写入 matched_orders、interval_rows、unshipped_orders、temporal_rows、orders_without_rate、usd_gap_krw。

参考

定义四个 Source 并统计行数

运行 flink-up 后,在 /root/flink/joins/ddl.sql 中定义 users · orders · shipments · rates(在 order_time · ship_time · update_time 上设置水位线)。用 sql-client.sh -i ddl.sql -f count.sql 运行以批处理模式、用 tbl、n 两列输出四张表行数的 /root/flink/joins/count.sql,并保存到 /root/flink/joins/count.out。

水位线写成 WATERMARK FOR 时间列 AS 时间列 的形式。每个文件都按各自的时间顺序排列,所以不设延迟也没有迟到行。把四个 SELECT 用 UNION ALL 连起来,就是一张表。

常规连接——不看时间,记住两侧

在 /root/flink/joins/regular.sql 中以流式模式把 orders 和 users 按 user_id 做内连接,输出按 tier 分组的 orders(条数)和 amount(amount 之和);并把输出保存到 /root/flink/joins/regular.out。

常规连接不需要时间条件,只用等值。u41–u44 是 users 中没有的用户,所以在内连接中被排除。请从结果中确认,下午注册的用户(signup_time 晚于订单)的订单是否也连上了——常规连接与到达顺序无关,会找出过去和未来的所有配对。

区间连接——只连上两小时内的发货

在 /root/flink/joins/interval.sql 中写区间连接:把 orders 和 shipments 按 order_id 连接,但只保留 ship_time 在订单时间起两小时内(含两端)的记录,并把 order_id、ship_id、delay_s(延迟秒数,TIMESTAMPDIFF(SECOND, ...))保存到 /root/flink/joins/interval.out。

区间连接必须有一个等值条件,以及把两侧时间绑在一起的范围。BETWEEN a AND b 包含两端。有延迟恰好为 0 秒和 7200 秒的发货,也有 7201 秒的发货。分两次发货的订单会变成两行。

外连接形式的区间连接——没有按时配对的订单

在 /root/flink/joins/unshipped.sql 中用 orders LEFT JOIN shipments 写出查询,输出两小时内一次都没有发货的订单的 order_id、order_time,并把输出保存到 /root/flink/joins/unshipped.out。

时间条件放在 ON 子句里,再用 WHERE 只选出因没有配对而用 null 填充的行。不仅要包含完全没有发货记录的订单,还要包含超过两小时才发出的订单。null 行会在水位线越过(订单时间 + 2 小时)之后输出一次。

时态连接——订单时刻的汇率

在 /root/flink/joins/temporal.sql 中,用 rates 创建按币种只保留最新行的版本视图 rates_v,用 FOR SYSTEM_TIME AS OF o.order_time 连接 orders,输出 order_id、currency、rate、amount_krw(amount × rate)。输出保存到 /root/flink/joins/temporal.out。

rates 是 append-only 的,不能声明主键。建一个只保留 ROW_NUMBER() OVER (PARTITION BY currency ORDER BY update_time DESC) 为 1 的行的视图,它就是以 currency 为主键、以 update_time 为事件时间的版本视图。如果是内连接,早于第一个汇率的订单会被排除。

同一个视图用常规连接——过去的订单会被重新计算

在 /root/flink/joins/latest.sql 中,把第 5 步的 rates_v 不加 FOR SYSTEM_TIME AS OF,按 currency 与 orders 连接,输出各币种的 orders(条数)和 total_krw(amount × rate 之和),并把输出保存到 /root/flink/joins/latest.out。

版本视图在汇率每次变化时都会输出更新。常规连接接收这些更新后,会把过去的订单也重新计算,所以输出中能看到 -U。最终合计与时态连接的合计不同——想一想它是按哪个汇率算出来的。

从三种连接的执行计划中读取状态

把第 2、3、5 步的连接改写成写入 blackhole Sink 的 INSERT,用 COMPILE PLAN 生成 /root/flink/joins/regular-plan.json、/root/flink/joins/interval-plan.json、/root/flink/joins/temporal-plan.json。(常规连接为 orders⋈users,区间连接为 orders⋈shipments,时态连接为 orders⋈rates_v)

形式是 COMPILE PLAN 'file:///路径.json' FOR INSERT INTO Sink SELECT ...。文件已存在会报错,所以重新导出时要先删除再运行。生成后用 jq 浏览 nodes 的 type 和 state——要点在于哪个连接节点有 leftState · rightState,哪个节点没有。

报告——每种连接连上了什么、漏掉了什么

在 /root/flink/joins/report.json 中写入 matched_orders(用常规连接与用户连上的订单数)、interval_rows(区间连接结果的行数)、unshipped_orders(第 4 步的行数)、temporal_rows(时态连接结果的行数)、orders_without_rate(全部订单 − temporal_rows)、usd_gap_krw(latest.out 中 USD 的 total_krw − temporal.out 中 USD 的 amount_krw 之和,小数)。

全都可以从前面保存的输出中抄过来。带有变更日志的表,最后一个 +I/+U 就是最终值。用 grep 只选出“| +I |”这样的行,再用 awk -F'|' 切分列,就能数出来。整数列写整数,usd_gap_krw 写数字。