用三种连接关联订单
目标
对同一个订单流,分别用常规连接、区间连接和事件时间时态连接连上用户、发货和汇率,并通过输出和执行计划确认各种连接记住了什么、是否会修正结果。
为什么重要
流式连接不知道配对何时到来,所以要把对侧的行堆在状态里。没有时间条件就永远堆,给出时间范围就堆到水位线越过为止,如果是版本表,则只持有需要的版本。选错连接类型,状态会无限增长,或者过去的结果会被悄悄改变。本实验的评分器不会查询集群,而是把你保存的 sql-client 输出和计划 JSON,与直接从源 CSV 计算出的值进行核对。
步骤
- 用
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)。 - 在 /root/flink/joins/regular.sql 中写一个流式查询:把 orders 和 users 按
user_id连接,输出各等级的orders(订单数)和amount(合计),并把输出保存到 /root/flink/joins/regular.out。 - 在 /root/flink/joins/interval.sql 中写区间连接,连上从订单时间起两小时内(含两端)发出的货,把
order_id、ship_id、delay_s(秒)保存到 /root/flink/joins/interval.out。 - 在 /root/flink/joins/unshipped.sql 中写外连接形式的区间连接,输出两小时内没有发货的订单的
order_id、order_time,并把输出保存到 /root/flink/joins/unshipped.out。 - 用 rates 创建版本视图
rates_v,用时态连接连上订单时刻的汇率,把输出order_id、currency、rate、amount_krw(amount × rate)的 /root/flink/joins/temporal.sql 的输出保存到 /root/flink/joins/temporal.out。 - 把同一个
rates_v不加FOR SYSTEM_TIME AS OF连接,把输出各币种orders和total_krw的 /root/flink/joins/latest.sql 的输出保存到 /root/flink/joins/latest.out。 - 用
COMPILE PLAN把三种连接(第 2、3、5 步)的执行计划导出为 /root/flink/joins/regular-plan.json、/root/flink/joins/interval-plan.json、/root/flink/joins/temporal-plan.json。 - 在 /root/flink/joins/report.json 中写入
matched_orders、interval_rows、unshipped_orders、temporal_rows、orders_without_rate、usd_gap_krw。
参考
- 源数据(没有表头的 CSV,时间以秒为单位,每个文件按各自的时间顺序排列):
joins_users.csv=user_id, tier, signup_time·joins_orders.csv=order_id, user_id, currency, amount, order_time·joins_shipments.csv=ship_id, order_id, ship_time·joins_rates.csv=currency, rate, update_time(rate 以韩元计,请按DECIMAL(10, 4)读取)。都在/opt/lab/fixtures/data/中。 - 只想写一次定义,可以用
sql-client.sh -i ddl.sql -f 쿼리.sql > 쿼리.out 2>&1(占位符均为查询文件名)——-i文件中的 CREATE 语句会先执行。 - 流式结果最前面会多出
op列(+I · -U · +U · -D)。作业在文件末尾结束,此时水位线会推进到底,剩余的区间和版本都会被处理。 COMPILE PLAN接受INSERT INTO语句。请创建一个'connector' = 'blackhole'的 Sink 来用。同一路径下已有文件时不会覆盖而是报错,所以重新导出前要先删除。- 常见错误:外连接形式的区间连接,如果把时间条件放在
WHERE里,null 行会被过滤掉。时态连接的右侧必须有主键,而 append-only 的 Source 必须用去重视图来做。 - 官方文档:Joins · Versioned Tables · Deduplication · SQL Client
定义四个 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 写数字。