Apache Spark — 慢作业的答案在执行计划和事件日志里
用 SQL 写还是用 DataFrame 写都会得到相同的计划,窗口函数不会减少行数
一句话总结
SQL 字符串和 DataFrame API 是同一个引擎转换成同一个计划的两个入口。选择的标准不是哪个更快,而是哪个更易读。聚合是把行归并、减少行数,而窗口函数则保留原来的行,同时查看相邻的行——而且窗口的结果取决于帧是怎么设定的。
为什么有两个入口
数据团队里,有人用 SQL 思考,有人用 Python 思考。于是很快就会有这样的说法流传:“用 SQL 写更快”,或者“DataFrame 的优化做得更好”。两种说法都不对。
Spark SQL 指南的第一段就是答案。Spark SQL 比 RDD 掌握更多关于数据和计算的结构信息,并利用这些信息做更多优化。而且在计算结果时,无论用哪种 API 或语言来表达,使用的都是同一个执行引擎。文档写道,正因为有这种统一,每一步转换都可以选择更自然的 API 来回切换。
“同一个引擎”这句话是具体的。SQL 字符串被解析成逻辑计划,DataFrame 方法调用也是直接堆叠出逻辑计划。从这里开始,它们经过同一个分析器、同一个优化器、同一个物理计划器。所以把同一个问题用两种方式写出来,再打印 explain(),物理计划是一样的。本实验中要亲手比较这一点。
工作原理
要用 SQL 调用 DataFrame,需要一个名字。入门文档展示了用 createOrReplaceTempView 把 DataFrame 注册成临时视图,再用 spark.sql 查询的样子。同一页写道,临时视图的作用域是会话,创建它的会话结束时就会消失。视图不会复制数据。它只是贴在计划上的一个名字。
orders.createOrReplaceTempView("orders")
by_sql = spark.sql("""
SELECT date_trunc('month', order_ts) AS month, sum(qty) AS units
FROM orders WHERE status = 'paid' GROUP BY 1""")
by_api = (orders.where(F.col("status") == "paid")
.groupBy(F.date_trunc("month", "order_ts").alias("month"))
.agg(F.sum("qty").alias("units")))
聚合的物理计划几乎总是同一个形状。看 EXPLAIN 文档的示例,自下而上的顺序是 HashAggregate(... partial_sum ...) → Exchange hashpartitioning(k, 200) → HashAggregate(... sum ...)。每个分区先把自己那份归并成部分和,再按键做 shuffle,然后把汇集起来的部分和再相加。通过 shuffle 发送的不是原来的行,而是每个键归并成一行的部分和。所以像求和、计数这样的聚合,即使数据很大,shuffle 也比想象中小。示例里的 200 是 shuffle 分区数,在讲 shuffle 的模块里会另外讲到。
窗口函数——不减少行数的聚合
groupBy 每个类别只留下一行。“每个类别销量前 3 的商品”就没法这样解决。要让商品的行原封不动,同时在旁边的一列里写上它在同一类别里排第几。这就是窗口函数。SQL 参考中的窗口函数部分把它们分为排名函数(RANK、DENSE_RANK、ROW_NUMBER 等)和分析函数(LAG、LEAD、FIRST_VALUE 等),并定义了用 ROWS 或 RANGE 来写帧的语法。
排名函数的差别出现在并列的时候。如 rank 文档所说,dense_rank 在并列之后不会留下空缺的名次,而 rank 如果有三个并列第 2 名,下一名就记为第 5 名。如果用 rank <= 3 来筛选“前 3”,由于并列,可能会出现四行以上;如果用 row_number 来筛选,并列者谁能入选就交给排序顺序决定。哪种是对的,由业务决定。
更隐蔽的陷阱是默认帧。Window 文档写道,没有排序时,默认使用整个分区(按行,从头到尾);有排序时,默认使用按范围、从开头到当前行逐渐增长的帧。根据 rangeBetween 文档,范围边界以 ORDER BY 的值为准,而不是行的位置。所以按日期排序求累计和时,同一日期的行会互相把对方算进“到当前行为止”,同时得到相同的累计值。想要一行一行增长的累计和,就必须用 rowsBetween(Window.unboundedPreceding, Window.currentRow) 写出按行的帧。
精确的数和近似的数
要精确统计去重客户数,必须把所有客户 ID 汇集起来去重,所以值本身就要做 shuffle。approx_count_distinct 接收允许的相对标准差 rsd 来估算,默认值是 0.05。文档写道,如果 rsd 需要小于 0.01,那还不如用 count_distinct 更高效。内置聚合函数列表说明这种估算是通过 HyperLogLog++ 完成的。用在像仪表板里的每日访客这种 5% 误差无所谓的地方,而不要用在像结算这样一个人都不能错的地方。
在现场相遇的样子
SQL 字符串里的错误是在运行时才出现的。列名写错了,Python 编辑器一声不吭,直到调用 spark.sql 的那一刻才出现分析错误。如果在 DataFrame API 里把列写成字符串,也是一样。无论哪种,都会在行动算子之前的分析阶段被发现,所以测试最便宜的办法,是用小数据完整地跑一遍。
两种方式混着用是常态。复杂的窗口和连接用 SQL 写,重复生成列这类事情则用 Python 循环和 DataFrame API 来写。计划是一样的,所以混用也没有损失。
临时视图在会话之外是看不到的。在别的应用或别的会话里用同样的名字去调用,会找不到表。要在应用之间传递数据,必须写成文件或表。
实际工作中真正重要的事
- SQL 和 DataFrame 是同一个引擎、同一个计划。按是否易读而不是性能来选,有怀疑时用 explain 比较。
- 聚合是部分聚合 → shuffle → 最终聚合。被 shuffle 的不是原来的行,而是缩减后的部分结果。
- 窗口不会减少行数。前 N 用排名函数来排,并按业务选择并列的规则。
- 累计和要明确写出帧。有排序的窗口默认是范围帧,所以值相同的行会被一起算进去。
- 近似去重计数只用在误差无所谓的地方。默认 rsd 是 0.05。
下一项实验要做什么
把订单数据注册成临时视图,用 SQL 得出按月销售额,再用 DataFrame API 重写同样的结果,并用 explain 确认两边的计划相同。用窗口函数挑出每个类别销量前 3 的商品,用按行的帧计算累计和,用 lag 计算环比上月的变化。最后,分别用近似函数和精确方法统计去重客户数,测出误差,并整理成报告。