2.3 SQL查询实战:声明式入口的巡检


2.3 SQL 查询实战:声明式入口的巡检

本节摘要:Spark SQL 允许直接执行标准 SQL,且 SQL 与 DataFrame 在引擎内汇入同一执行计划。本节演示建视图、临时表、混编写法与窗口函数实战,并从执行视角解释"SQL 只是外衣"——同一段逻辑两种写法,explain 输出可以逐字相同。

从建视图开始

from pyspark.sql import SparkSession from pyspark.sql import functions as F spark = SparkSession.builder.appName("SqlLab").master("local[*]").getOrCreate() orders = spark.read.option("header", True).csv("orders") users = spark.read.option("header", True).csv("users") orders.createOrReplaceTempView("orders") # 视图:只是给逻辑计划起名,不物化 users.createOrReplaceTempView("users")

巡检要点:createOrReplaceTempView 注册的是逻辑计划的一个名字,引擎没有复制任何数据。数据物化只发生在缓存或写盘时。全局临时视图(global_temp 前缀)跨会话可见,普通临时视图只在当前会话有效。

同一逻辑的两种外衣

sql_result = spark.sql(""" SELECT u.region, count(*) AS order_cnt, sum(o.amount) AS gmv FROM orders o JOIN users u ON o.user_id = u.user_id WHERE o.amount > 50 GROUP BY u.region HAVING count(*) > 10 ORDER BY gmv DESC """) df_result = (orders.join(users, "user_id") .filter("amount > 50") .groupBy("region") .agg(F.count("*").alias("order_cnt"), F.sum("amount").alias("gmv")) .filter("order_cnt > 10") .orderBy(F.desc("gmv"))) sql_result.explain() == df_result.explain() # 物理计划文本一致(格式忽略时)

两种写法经 Catalyst 处理后落在同一物理计划上。所以选哪种纯看团队习惯:分析同学写 SQL,工程同学写 DataFrame,混在同一作业里也没有额外开销。

窗口函数:一次 Shuffle 上的多层计算

from pyspark.sql import Window w_rank = Window.partitionBy("region").orderBy(F.desc("amount")) w_cum = Window.partitionBy("region").orderBy("pay_time").rowsBetween(Window.unboundedPreceding, 0) (report := (orders .filter("amount > 50") .withColumn("region_rank", F.rank().over(w_rank)) .withColumn("cum_amount", F.sum("amount").over(w_cum)) .filter(F.col("region_rank") <= 3))) report.show()

引擎视角的关键:两个窗口若分区与排序键相同,Catalyst 能把它们合并到同一次 Shuffle 后的窗口运算里;分区键不同的窗口则各付一次 Shuffle。写多条窗口列时对齐分区键,是 SQL 调优里性价比最高的一招。

执行视角的写法守则

  • 谓词尽早写:WHEREJOIN 前与后,Catalyst 多数时候能帮你下推,但显式先过滤是习惯保险
  • SELECT 只点名要用的列,列裁剪对列存格式直接减少 IO
  • DISTINCTGROUP BY 全列等价,代价都是一次全列 Shuffle
  • 大表之间 Join 前先聚合瘦身,比事后过滤便宜得多

💡 关键直觉:SQL 的"声明式"意味着你交出的是意图,引擎交回的是计划。巡检慢 SQL 的正确姿势永远是先 explain、再谈改写——凭直觉改 SQL,常常是在优化 Catalyst 早已处理过的东西。

巡检案例:一条漏斗报表 SQL 的三轮瘦身

背景:运营要"近 30 天各区 Top3 商品的复购漏斗",源表订单明细一亿两千万行。第一版 SQL 由报表同学照业务口径直译:三张表全字段 join、join 完再算窗口、窗口完再过滤商品名。操作:先 explain 观察再逐轮改写。第一轮只做两处:SELECT 点名列替代星号、把日期过滤从外层提到 join 之前。结果:扫描量降四成——列裁剪让 Parquet 只读六列,日期谓词触发了分区裁剪,整月目录直接跳过二十九个。第二轮处理窗口:原写法先 join 出宽表再开两个分区键不同的窗口,两轮 Shuffle;改成先按商品预聚合出候选集(缩到二十万行)再 join 维表开窗口。结果:窗口阶段的输入从一亿行降到二十万行,总耗时从 41 分钟压到 7 分钟。第三轮核对口径:把 rank 换成 dense_rank 消除并列商品挤占名次的争议,属于业务正确性修正而非性能项。解读:瘦身三步的次序有讲究——先砍输入(裁剪与过滤),再砍宽依赖的输入规模(先聚合后窗口),最后才动窗口内部细节;顺序颠倒会在错误层面白费力气。变式:若源表按日分区且查询固定跨整月,可进一步用增量物化的中间表替代全月重算,报表场景的终局往往是调度问题而不再是引擎问题。

-- 第二轮的核心改写示意:先聚合出候选集再开窗口 SELECT region, goods_name, replay_cnt FROM ( SELECT region, g.goods_name, count(DISTINCT repeat_uid) AS replay_cnt, dense_rank() OVER (PARTITION BY region ORDER BY count(DISTINCT repeat_uid) DESC) AS rk FROM agg_candidate c JOIN goods_dim g ON c.goods_id = g.goods_id GROUP BY region, g.goods_name ) t WHERE rk <= 3

本节要点回顾

  • 视图不物化:临时视图只是逻辑计划的别名,缓存才物化
  • 殊途同归:SQL 与 DataFrame 汇入同一物理计划,混用零额外代价
  • 窗口合并规则:分区与排序键一致的窗口共享一次 Shuffle
  • explain 先行:改写 SQL 前先看计划,避免无效优化

作者与出处
原作者: 灏天文库
来源:灏天文库
整理: 灏天文库整理
由灏天文库平台收录,内容或由平台用户上传,仅供学习交流
发布者: 作者: 灏天文库 转发
评论区 (0)
U