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

Apache Spark — 慢作业的答案在执行计划和事件日志里

用 SQL 和 DataFrame 计算月度销售额、分类排名和累计和

在 TT Lab 中继续学习

目标

把同一个问题分别用 Spark SQL 和 DataFrame API 解决,确认得到相同的物理计划,并用窗口函数计算排名、累计和、环比上月。最后比较精确去重计数和近似去重计数。

为什么重要

在 Spark 里,SQL 和 DataFrame 不是两个引擎。两者都被转换成同一个逻辑计划,经过同一个优化器(Catalyst),再以同一个物理计划执行。所以“SQL 更快”或“DataFrame 更快”,多半是个错误的问题。无论团队用什么来写,只要看计划,马上就知道相同还是不同。 聚合之后最常用的就是窗口函数。它能在不减少行数的情况下查看相邻的行,排名、累计和、环比上月这类报告里的数字,全都出自这里。窗口的框架(分区、排序、范围)如果设错,不会报错,却会得出错误的数字,所以必须能用话把这个框架讲清楚。 统计去重数是 shuffle 很大的运算。只要接受百分之几的误差,用 HyperLogLog++ 就能便宜得多地统计。仪表板用近似,账单用精确——用哪个,取决于这个数字的用途。

步骤

  1. 在 /root/spk/sql/common.py 中实现 load(spark),按 schema 读取订单和商品,并注册成临时视图 orders 和 products;在 /root/spk/sql/sql_monthly.py(应用 spk-sql-monthly)中用 SQL 把按月销售额(列 month、revenue)以带表头行的 CSV 写入 /root/spk/sql/out/monthly_sql。
  2. 在 /root/spk/sql/df_monthly.py(应用 spk-sql-df)中不用 spark.sql,而用 DataFrame API 把同样的结果写入 /root/spk/sql/out/monthly_df。
  3. 在 /root/spk/sql/plans.py(应用 spk-sql-plans)中,把两种方式的 explain(mode="formatted") 输出分别保存到 /root/spk/sql/out/plan_sql.txt 和 /root/spk/sql/out/plan_df.txt,然后把两个 DataFrame 各自用 collect() 真正运行一遍。
  4. 在 /root/spk/sql/top3.py(应用 spk-sql-top3)中,按类别(category)把销售额前 3 的商品,以 category、product_id、revenue、rank 这几列写入 /root/spk/sql/out/top3。如果并列,product_id 靠前的排在上面。
  5. 在 /root/spk/sql/running.py(应用 spk-sql-running)中,把 2026 年 1 月各渠道的日销售额和渠道内的累计和,以 channel、day、revenue、running 这几列写入 /root/spk/sql/out/running。
  6. 在 /root/spk/sql/mom.py(应用 spk-sql-mom)中,把各渠道的月销售额、上月销售额、增长率(百分比,四舍五入到小数点后第二位),以 channel、month、revenue、prev_revenue、growth_pct 这几列写入 /root/spk/sql/out/mom。第一个月的上月值留空。
  7. 在 /root/spk/sql/distinct.py(应用 spk-sql-distinct)中,精确地(countDistinct)以及用 approx_count_distinct(rsd=0.05),统计已完成支付订单的下单客户数,并以 {"exact": 정수, "approx": 정수}(占位符为整数)写入 /root/spk/sql/out/distinct.json。
  8. 在 /root/spk/sql/report.md 中以 ## SQL 과 DataFrame(韩文,意为“SQL 与 DataFrame”)、## 윈도 함수(韩文,意为“窗口函数”)、## 근사 집계(韩文,意为“近似聚合”)三节写成报告。第三节放入第 7 步的两个数字。

参考

用临时视图和 SQL 得出按月销售额

在 /root/spk/sql/common.py 中实现 load(spark),按 schema 读取订单和商品,并注册成临时视图 orders 和 products;再以应用名 spk-sql-monthly 创建 /root/spk/sql/sql_monthly.py,用 spark.sql 把已完成支付订单的按月销售额(列 month 为 yyyy-MM,revenue 为 qty×price 之和)以带表头行的 CSV 写入 /root/spk/sql/out/monthly_sql。

临时视图是只在这个 SparkSession 里才看得到的名字。注册视图不会读取任何东西——它只是给计划贴上一个名字。用 date_format(order_ts, 'yyyy-MM') 得到月份,再与商品连接并乘以价格。

用 DataFrame API 解决同一个问题

以应用名 spk-sql-df 创建 /root/spk/sql/df_monthly.py,不使用 spark.sql,而是用 DataFrame API(where、join、groupBy、agg)把与第 1 步相同的结果,以带表头行的 CSV(month、revenue)写入 /root/spk/sql/out/monthly_df。

用 F.date_format("order_ts", "yyyy-MM").alias("month") 来分组,用 F.sum(F.col("qty") * F.col("price")) 来求和。连接键的名字相同时,像 join(p, "product_id") 这样以字符串传入,让该列只保留一个。

看看两种方式的物理计划是否相同

以应用名 spk-sql-plans 创建 /root/spk/sql/plans.py,分别构造第 1、2 步的两个查询,把 explain(mode="formatted") 的输出分别保存到 /root/spk/sql/out/plan_sql.txt 和 /root/spk/sql/out/plan_df.txt,然后把两个 DataFrame 各自用 collect() 真正运行一遍。

formatted 模式在上面打印算子树,下面打印带编号的算子说明。列编号(像 #12 这样的)每个查询都会不同,但算子的名字和顺序必须相同。真正运行之后,事件日志里的 SQL 执行记录中会留下同样的计划文档——评分器会检查你的文件是否与该记录相同,以及两个计划的算子是否相同。

每个类别销量前 3 的商品——窗口排名

以应用名 spk-sql-top3 创建 /root/spk/sql/top3.py,求出按类别(category)和商品划分的已完成支付销售额,对每个类别按销售额降序(并列时按 product_id 升序)取第 1–3 名,以带表头行的 CSV(category、product_id、revenue、rank)写入 /root/spk/sql/out/top3。

窗口是 Window.partitionBy("category").orderBy(...)。rank() 给并列的值相同的编号,并跳过下一个编号,所以如果想要恰好 3 个,就要用 row_number(),并把并列规则放进排序里。

各渠道累计和——确定窗口的范围

以应用名 spk-sql-running 创建 /root/spk/sql/running.py,把 2026 年 1 月已完成支付订单的各渠道日销售额(day 是日期),以及在渠道内按日期顺序累加的累计和 running,以带表头行的 CSV(channel、day、revenue、running)写入 /root/spk/sql/out/running。

累计和是在 Window.partitionBy("channel").orderBy("day").rowsBetween(Window.unboundedPreceding, Window.currentRow) 之上的 sum。要先按日期汇总成一行,再套上窗口——如果在汇总之前套,同一天的每个订单的累计和都会不同。

环比上月——用 lag 看相邻的行

以应用名 spk-sql-mom 创建 /root/spk/sql/mom.py,把各渠道的月销售额、用 lag 取来的上月销售额 prev_revenue、增长率 growth_pct(=(本月-上月)/上月×100,四舍五入到小数点后第二位),以带表头行的 CSV(channel、month、revenue、prev_revenue、growth_pct)写入 /root/spk/sql/out/mom。渠道的第一个月,上月和增长率必须为空。

F.lag("revenue").over(Window.partitionBy("channel").orderBy("month")) 会取来同一渠道紧邻的上一个月的值。第一个月没有前一行,所以是 null;含有 null 的除法结果也是 null,所以不需要另外处理。

精确去重计数和近似去重计数

以应用名 spk-sql-distinct 创建 /root/spk/sql/distinct.py,用 countDistinct 和 approx_count_distinct(rsd=0.05) 一次性统计已完成支付订单的下单客户数,并以 {"exact": 정수, "approx": 정수}(占位符为整数)写入 /root/spk/sql/out/distinct.json。

近似的一边是合并 HyperLogLog++ 草图。不需要通过 shuffle 来汇集去重的值,所以数据越大越便宜。rsd 是相对标准误差的目标值,所以要亲自算一算结果与精确值相差百分之几。评分器还会检查事件日志的计划里是否真的包含近似函数。

留下相同的计划、窗口的框架、近似的误差

在 /root/spk/sql/report.md 中写出 ## SQL 과 DataFrame(韩文,意为“SQL 与 DataFrame”)、## 윈도 함수(韩文,意为“窗口函数”)、## 근사 집계(韩文,意为“近似聚合”)三节。第三节用数字放入第 7 步的精确值和近似值。

第一节写下两个计划中相同的几个算子名,第二节用一行分别写出排名、累计和、环比上月里窗口是怎么设定的。第三节写上两个数字和它们之差的百分比就行。