前三篇用 DataFrame API 完成了从读取到进阶分析的全流程。但团队里不是所有人都写 Python——数据分析师习惯 SQL,BI 工具也只认 SQL。Spark SQL 让同一份数据既能用 DataFrame 操作,也能用标准 SQL 查询,两种写法底层走同一个优化器,性能一致。
本文要点
- spark.sql() 的基本用法:能查什么、返回什么
- 临时视图(Temp View)与全局临时视图(Global Temp View)的区别与使用场景
- DataFrame API 与 SQL 的互转:什么时候用哪种写法
- 常用内置函数分类速查(字符串、日期、数学、聚合、窗口)
- UDF 原理:为什么 Python UDF 慢、Pandas UDF 为什么快、什么时候不得不用 UDF
- 贯穿案例:从纯 DataFrame 到 SQL + 临时视图,逐步优化城市销售日报
阅读前提醒:Spark SQL 和你熟悉的 MySQL / PostgreSQL 不完全一样。它兼容标准 SQL 语法的同时,还支持很多大数据特有的能力——比如 STRUCT、ARRAY、MAP 等复杂类型、窗口函数的扩展用法、以及直接查询文件的语法。看到陌生的语法不用慌,后面会专门讲。
零 为什么需要 Spark SQL
1 一个现实问题
数据团队通常有两类人:
| 数据开发 | Python / Scala | 写 ETL 管道、复杂逻辑 |
| 数据分析师 | SQL | 快速取数、临时查询、报表 |
如果只有 DataFrame API,分析师必须学 Python 才能取数;如果只有 SQL,开发写复杂管道时表达力不够。Spark SQL 的解法是:同一份数据,两种接口。
# 开发用 DataFrame API 写 ETL
df = spark.table("dw.orders").filter(col("amount") > 0)
# 分析师用 SQL 查同一张表
result = spark.sql("SELECT city, COUNT(*) FROM dw.users GROUP BY city")
两种写法底层都经过 Catalyst 优化器,生成相同的物理执行计划。不存在"SQL 比 DataFrame 快"或反过来的情况——性能一致,选择取决于表达习惯和场景。
2 先建立一个整体认知
在开始之前,先搭好一个框架,后面的内容就不容易乱。
Spark 内核只有一个(Catalyst 优化器 + RDD 执行引擎),但对外开了两扇门:
| 入口 | df.select(), df.filter() … | spark.sql("SELECT …") |
| 内置函数 | pyspark.sql.functions 里的 count()、sum()、when() … | COUNT, SUM, CASE WHEN … |
| 自定义函数 | 用 udf() / pandas_udf() 包装 Python 函数 | 用 spark.udf.register() 注册 Python 函数 |
| 底层 | 都走 Catalyst 优化器 → 物理计划 → RDD 执行 | 同左 |
每扇门旁边都有两个工具箱:
- 内置工具箱——Spark 自带的几百个函数,两门通用,只是写法不同
- 自定义工具箱——你自己写的 Python 逻辑,通过 UDF 这个"转换器",塞进两边的工具箱里
Spark 内核(Catalyst + RDD)
┌──────────────────────────┐
│ │
DataFrame API 门 SQL 门
┌───────────┐ ┌───────┐
│ 内置函数 │ │ 内置 │
│ count/sum │ │COUNT/ │
│ when/col │ ←── UDF ──→ │SUM/…│
│ … │ (转换器) │ │
├───────────┤ ├───────┤
│ 自定义 │ │ 自定义 │
│ udf()包装 │ │register│
│ 的Python │ │ 的函数 │
│ 函数 │ │ │
└───────────┘ └───────┘
一句话总结:两个交互界面(DataFrame API / SQL),每个界面下又分内置的和自定义的。UDF 就是把你写的 Python 逻辑,转换成两个界面都能调用的形式。
3 本篇的主线任务
本篇围绕一个具体任务展开:做一张城市维度的销售日报。
需要输出的报表长这样:
| 上海 | VIP | 1,200 | 580,000 | 483.33 |
| 上海 | 普通 | 8,500 | 1,200,000 | 141.18 |
| 北京 | VIP | 980 | 450,000 | 459.18 |
| … | … | … | … | … |
涉及三张表:
- dw.users — 用户表(user_id, city, age, …)
- dw.orders — 订单表(order_id, user_id, product_id, amount, order_time, order_status, …)
- dw.products — 商品表(product_id, category, price, …)
我们会从纯 DataFrame 写法出发,逐步引入 SQL 和临时视图,看看每种写法各自的优劣。
一 先试试:纯 DataFrame 写日报
先用前几篇学的 DataFrame API 来写这个任务,看看感觉怎么样。
from pyspark.sql import SparkSession
from pyspark.sql.functions import col, countDistinct, sum, avg, when
spark = SparkSession.builder.appName("sql-demo").getOrCreate()
# 读取数据
users = spark.table("dw.users")
orders = spark.table("dw.orders")
products = spark.table("dw.products")
# 城市销售日报
daily_report = (
orders
# 1. 过滤有效订单
.filter(col("order_status") == "completed")
.filter(col("order_date") == "2026-08-21")
.filter(col("amount") > 0)
# 2. 关联用户表,拿到城市
.join(users, "user_id")
# 3. 关联商品表(虽然当前指标用不上,但业务方说下周要加品类维度)
.join(products, "product_id")
# 4. 派生用户分层
.withColumn("user_tier", when(col("amount") > 1000, "VIP").otherwise("普通"))
# 5. 按城市 + 分层聚合
.groupBy("city", "user_tier")
.agg(
countDistinct("user_id").alias("active_users"),
sum("amount").alias("gmv"),
avg("amount").alias("avg_order_value")
)
# 6. 按 GMV 降序
.orderBy(col("gmv").desc())
)
daily_report.show()
看起来还行?那再加几个需求试试:
加了窗口函数和多层聚合之后,DataFrame 的写法就开始有点绕了:
from pyspark.sql.window import Window
from pyspark.sql.functions import spark_sum, lag, round as spark_round
window_city = Window.partitionBy("city").orderBy(col("gmv").desc())
window_date = Window.partitionBy("city", "user_tier").orderBy("order_date")
# 加排名和环比
daily_report_with_metrics = (
daily_report
.withColumn("city_rank", row_number().over(window_city))
.withColumn("prev_day_gmv", lag("gmv", 1).over(window_date))
.withColumn(
"mom_growth_rate",
spark_round(
(col("gmv") – col("prev_day_gmv")) / col("prev_day_gmv"),
4
)
)
)
逻辑越复杂,DataFrame 的链式写法就越像"叠罗汉"——每加一列就多一行 withColumn,窗口定义散落在各处,读起来要来回跳。
那换 SQL 写呢?SQL 的 WITH 子句(CTE)可以把逻辑分层,窗口函数写在 SELECT 里,读起来更像自然语言。
但问题来了:上面的 daily_report 是个 DataFrame,我能直接用 SQL 查它吗? 别急,先从最简单的 SQL 用法开始。
二 Spark SQL 初体验:spark.sql() 直接查
2.1 最基本的用法
spark.sql() 是 Spark SQL 的入口,不需要任何额外配置,拿到 SparkSession 就能直接写:
# 直接查 Paimon 表
result = spark.sql("SELECT city, COUNT(*) as user_count FROM dw.users GROUP BY city")
result.show()
spark.sql() 的返回值是 DataFrame——这一点非常重要。SQL 查询的结果天然就是 DataFrame,可以继续用 DataFrame API 操作:
# 先写 SQL 查,再用 DataFrame API 继续处理
result = spark.sql("SELECT city, COUNT(*) as user_count FROM dw.users GROUP BY city")
result.filter(col("user_count") > 1000).orderBy(col("user_count").desc()).show()
2.2 SQL 能查询的三类数据源
SQL 语句中 FROM 后面可以跟三类对象:
| Catalog 表 | dw.users | Paimon / Hive 表,元数据持久化存储,重启不丢 |
| 文件路径 | parquet./path/to/data`` | 直接读文件,不需要提前注册 |
| 视图 | users | 临时注册的逻辑表(后面会讲) |
# 1. 查 Paimon 表
spark.sql("SELECT * FROM dw.users").show()
# 2. 直接查文件
spark.sql("SELECT * FROM parquet.`/data/orders_parquet`").show()
spark.sql("SELECT * FROM csv.`/data/users.csv` WHERE _c0 = '1001'").show()
# 3. 查视图(需要先注册,下一节讲)
# spark.sql("SELECT * FROM user_view")
2.3 一个更实际的问题
纯 SQL 从零写日报当然可以,但实际工作中更常见的情况是:前面的数据清洗和关联已经有人用 DataFrame 写好了,我不想再从头写一遍 SQL。
比如团队里的数据开发已经写好了这段逻辑:
# 数据开发写的 ETL 中间结果
clean_orders = (
spark.table("dw.orders")
.filter(col("order_status") == "completed")
.filter(col("order_date") == "2026-08-21")
.filter(col("amount") > 0)
.join(spark.table("dw.users"), "user_id")
.join(spark.table("dw.products"), "product_id")
.withColumn("user_tier", when(col("amount") > 1000, "VIP").otherwise("普通"))
)
clean_orders 已经是清洗好、关联好的 DataFrame 了。我想直接基于它,用 SQL 做聚合和报表——因为 SQL 写聚合和窗口函数更顺手。
那问题来了:我能直接 spark.sql("SELECT * FROM clean_orders") 吗?
答案是:不能直接查,但可以通过临时视图来实现。下一节就讲这个。
三 临时视图:让 DataFrame 能被 SQL 查询
3.1 一个常见的新手误区
初学者常常会写这样的代码:
# ❌ 错误!SQL 引擎不认识 Python 变量
clean_orders = spark.table("dw.orders").filter(col("order_status") == "completed")
result = spark.sql("SELECT * FROM clean_orders") # 报错:Table or view not found: clean_orders
为什么不行?因为 SQL 引擎和 Python 是两个独立的运行时。SQL 引擎有自己的 Catalog(元数据目录),它只认识自己 Catalog 里注册过的"表",不认识 Python 里的变量。
那如果我有一个 DataFrame 中间结果,想用 SQL 查它,怎么办?
答案是:把 DataFrame 注册成临时视图。
3.2 什么是临时视图
临时视图的作用可以用一句话概括:给 DataFrame 起一个 SQL 能识别的表名,把它注册到 SQL 引擎的 Catalog 里。
可以把 createOrReplaceTempView 理解为"给 Python 变量挂个 SQL 门牌号":clean_orders 这个 DataFrame 本来只活在 Python 里,挂上门牌后 SQL 引擎就能找到并查询了。
# 用 DataFrame 做前面的清洗和关联
clean_orders = (
spark.table("dw.orders")
.filter(col("order_status") == "completed")
.filter(col("order_date") == "2026-08-21")
.filter(col("amount") > 0)
.join(spark.table("dw.users"), "user_id")
.join(spark.table("dw.products"), "product_id")
.withColumn("user_tier", when(col("amount") > 1000, "VIP").otherwise("普通"))
)
# 注册临时视图——给 DataFrame 挂个 SQL 门牌号
clean_orders.createOrReplaceTempView("clean_orders")
# 现在可以用 SQL 直接查这个 DataFrame 了
daily_report = spark.sql("""
SELECT
city,
user_tier,
COUNT(DISTINCT user_id) as active_users,
SUM(amount) as gmv,
AVG(amount) as avg_order_value
FROM clean_orders
GROUP BY city, user_tier
ORDER BY gmv DESC
""")
createOrReplaceTempView("clean_orders") 做了两件事:
关键点:临时视图不保存数据,它只是一个指向 DataFrame 的指针。视图存活期间,DataFrame 的 lineage(血缘关系)一直有效,查询时才会真正执行。
createTempView 在视图已存在时会报错,createOrReplaceTempView 会覆盖。日常开发用后者更方便。
3.3 为什么需要临时视图
回到我们的日报任务。有了临时视图之后,工作流就变成了:
读取数据 → DataFrame 清洗关联 → 注册临时视图 → SQL 聚合报表
前半段用 DataFrame 写(清洗、关联、派生字段,链式调用很顺手),后半段用 SQL 写(聚合、窗口函数、分层逻辑,SQL 可读性更好)。各取所长。
从交互角度看,临时视图是 DataFrame API 和 SQL 之间的一座桥。但为什么感觉是"单向"的?
DataFrame ──注册临时视图──▶ SQL 可查询
▲ │
└───────── spark.sql() ────────┘ (直接返回 DataFrame,不需要中介)
- DataFrame 想让 SQL 查 → 需要注册临时视图(给 Python 变量起个 SQL 能识别的表名)
- SQL 结果想继续用 DataFrame API → 不需要中介,spark.sql() 直接返回 DataFrame
为什么是"单向"的?
不是因为 SQL 更底层,也不是技术上有什么本质不对称。纯粹是因为我们站在 Python 这一边写代码:
- 你从 Python 调 spark.sql(),是你主动发起的调用,返回值自然回到你手里——就是 DataFrame,不需要桥。
- 但你想让 SQL 引擎引用 Python 里的一个变量?不行。SQL 引擎有自己的 Catalog(元数据目录),它只认自己目录里登记过的东西,Python 变量对它来说是"外面的世界"。所以得主动注册一下。
如果反过来,你站在纯 SQL 环境里(比如用 JDBC 连 Spark),那情况就倒过来了:SQL 想调 Python 逻辑需要注册 UDF,而 UDF 的返回值自然回到 SQL 结果里。
结论:谁主动发起调用,谁的那边就不需要"桥";被动的那一方,需要对方主动把东西注册到自己的目录里才能用。
3.4 临时视图 vs 全局临时视图
临时视图的作用域是当前 SparkSession。如果另一个 SparkSession(比如另一个 Notebook 或另一个应用)也想查这个视图,就会找不到。
# Session A
clean_orders.createOrReplaceTempView("clean_orders")
spark.sql("SELECT * FROM clean_orders") # 正常
# Session B(另一个 SparkSession)
spark2.sql("SELECT * FROM clean_orders") # 报错:Table or view not found
全局临时视图(Global Temp View)解决跨 Session 共享的问题:
# Session A:注册全局临时视图
clean_orders.createOrReplaceGlobalTempView("clean_orders_global")
# Session B:也能查到(注意要加 global_temp 前缀)
spark2.sql("SELECT * FROM global_temp.clean_orders_global")
| 注册方法 | createOrReplaceTempView() | createOrReplaceGlobalTempView() |
| 作用域 | 当前 SparkSession | 所有 SparkSession(同一 SparkContext) |
| 查询方式 | SELECT * FROM view_name | SELECT * FROM global_temp.view_name |
| 生命周期 | Session 结束即销毁 | 应用(SparkContext)停止才销毁 |
| 典型场景 | 单次查询、脚本内使用 | 多 Notebook 共享、交互式分析 |
注意:全局临时视图的名字不是 clean_orders_global,而是 global_temp.clean_orders_global。global_temp 是 Spark 内置的特殊数据库名,不可更改。
3.5 视图与 Paimon 表的关系
前面一直在用 spark.table("dw.users") 直接读 Paimon 表。Paimon 表是持久化的表,元数据存在 Catalog 里,数据存在文件系统里,跨 Session、跨应用都能访问。
临时视图是瞬时的,只存在于内存中的 Catalog 注册表里,应用重启就没了。
# Paimon 表:持久化,重启后还在
spark.sql("SELECT * FROM dw.users")
# 临时视图:瞬时,重启后需重新注册
clean_orders.createOrReplaceTempView("clean_orders")
spark.sql("SELECT * FROM clean_orders")
实际工作中,分析师通常直接查 Paimon 表(dw.users、dw.orders),不需要注册视图。临时视图更多用在开发写复杂管道时,把中间结果注册成视图,方便分段调试或切换写法。
四 SQL 与 DataFrame API 的互操作
4.1 双向流转:日报的完整工作流
把前面的内容串起来,SQL 和 DataFrame API 之间可以无缝切换:
# 第一步:DataFrame 做清洗和关联
clean_orders = (
spark.table("dw.orders")
.filter(col("order_status") == "completed")
.filter(col("order_date") == "2026-08-21")
.filter(col("amount") > 0)
.join(spark.table("dw.users"), "user_id")
.join(spark.table("dw.products"), "product_id")
.withColumn("user_tier", when(col("amount") > 1000, "VIP").otherwise("普通"))
)
# 第二步:注册临时视图
clean_orders.createOrReplaceTempView("clean_orders")
# 第三步:SQL 做聚合 + 窗口函数
city_report = spark.sql("""
SELECT
city,
user_tier,
COUNT(DISTINCT user_id) as active_users,
SUM(amount) as gmv,
AVG(amount) as avg_order_value,
ROW_NUMBER() OVER (PARTITION BY city ORDER BY SUM(amount) DESC) as city_tier_rank
FROM clean_orders
GROUP BY city, user_tier
ORDER BY gmv DESC
""")
# 第四步:SQL 返回的还是 DataFrame,可以继续操作
# 比如加一个过滤条件
top_cities = city_report.filter(col("gmv") > 100000)
# 第五步:还可以再注册成视图,供下一段 SQL 使用
top_cities.createOrReplaceTempView("top_cities")
spark.sql("SELECT * FROM top_cities WHERE city = '上海'").show()
整个流程就是:DataFrame → 注册视图 → SQL → DataFrame → 再注册 → 再 SQL,来回切换没有任何额外开销。
4.2 SQL 与 DataFrame API 的等价对照
同一个逻辑,两种写法。掌握对照关系,看到 SQL 能想到 DataFrame,反之亦然。
| 选列 | df.select("user_id", "city") | SELECT user_id, city FROM df |
| 过滤 | df.filter(col("age") >= 18) | SELECT * FROM df WHERE age >= 18 |
| 去重 | df.dropDuplicates(["user_id"]) | SELECT DISTINCT * FROM df(近似) |
| 聚合 | df.groupBy("city").agg(count("*")) | SELECT city, COUNT(*) FROM df GROUP BY city |
| 排序 | df.orderBy(col("amount").desc()) | SELECT * FROM df ORDER BY amount DESC |
| Join | df1.join(df2, "user_id") | SELECT * FROM df1 JOIN df2 ON df1.user_id = df2.user_id |
| 窗口 | row_number().over(Window.partitionBy("city").orderBy(desc("amount"))) | ROW_NUMBER() OVER (PARTITION BY city ORDER BY amount DESC) |
| 条件分支 | when(col("amount") > 100, "high").otherwise("low") | CASE WHEN amount > 100 THEN 'high' ELSE 'low' END |
4.3 什么时候用哪种写法:三个真实场景
两种写法性能一致,选择只看你在干什么、你是谁。
场景一:你是数据分析师,上班第一件事——查个数
产品经理飞过来一句:“昨天上海的 GMV 多少?VIP 用户占比多少?”
这时候你打开 Notebook,直接写 SQL 最快:
SELECT
SUM(amount) as gmv,
COUNT(DISTINCT CASE WHEN amount > 1000 THEN user_id END) * 1.0
/ COUNT(DISTINCT user_id) as vip_ratio
FROM dw.orders
WHERE order_date = '2026-08-21'
AND city = '上海'
AND order_status = 'completed'
为什么用 SQL? 临时取数、快速验证、写起来像说话一样自然。这种"写了就跑、跑了就看、看完就扔"的查询,SQL 是最顺手的工具。
场景二:你是数据开发,在写每日 ETL 管道
需求是:每天凌晨 2 点,把前一天的订单数据清洗、关联三张表、派生十几个指标、写入 Paimon 结果表。
这种你大概率用 DataFrame API 写:
def daily_etl(target_date):
# 动态传参
df = spark.table("dw.orders").filter(col("order_date") == target_date)
# 链式调用,步骤清晰
result = (
df
.filter(col("order_status") == "completed")
.join(spark.table("dw.users"), "user_id")
.join(spark.table("dw.products"), "product_id")
.withColumn("user_tier", when(col("amount") > 1000, "VIP").otherwise("普通"))
.withColumn("order_month", date_format(col("order_time"), "yyyy-MM"))
# 动态列:列名从配置里读
.select(*required_columns)
)
# 结果写回 Paimon
result.write.mode("overwrite").saveAsTable("dws.daily_order_stats")
为什么用 DataFrame API?
- 可以写函数、传参数(target_date),SQL 字符串拼起来容易出错
- 列名可以动态传(df.select(*required_columns)),SQL 拼列名很痛苦
- 便于单元测试、版本管理、和 Python 其他库配合(Pandas、sklearn)
场景三:开发调试时——两边混着用
这是最常见的"日常操作"。你在用 DataFrame 写管道,写到一半想验证一下中间结果对不对。
以前的做法:每一步都 .show() 一下,肉眼看数据对不对。但如果逻辑复杂,.show() 只能看前 20 行,很难验证。
现在的做法:写几步就注册一个临时视图,然后用 SQL 做各种角度的验证:
# 前半段用 DataFrame 写 ETL
clean_data = (
spark.table("dw.orders")
.filter(col("order_date") == "2026-08-21")
.filter(col("order_status") == "completed")
.join(spark.table("dw.users"), "user_id")
.withColumn("user_tier", when(col("amount") > 1000, "VIP").otherwise("普通"))
)
# 注册成视图,方便用 SQL 验证
clean_data.createOrReplaceTempView("clean_data")
# 验证 1:各城市订单量分布
spark.sql("""
SELECT city, COUNT(*) as cnt
FROM clean_data
GROUP BY city
ORDER BY cnt DESC
""").show()
# 验证 2:VIP 用户的平均客单价是不是确实更高
spark.sql("""
SELECT user_tier, AVG(amount) as avg_amount
FROM clean_data
GROUP BY user_tier
""").show()
# 验证 3:有没有订单金额为负的异常数据
spark.sql("""
SELECT COUNT(*) as negative_count
FROM clean_data
WHERE amount < 0
""").show()
# 验证完没问题,继续用 DataFrame 往下写
final_result = clean_data.groupBy("city", "user_tier").agg(...)
为什么混用? DataFrame 写管道(结构化、可复用),SQL 做验证(灵活、直观、想怎么查怎么查)。各干各擅长的事。
一句话总结:临时取数用 SQL,写管道用 DataFrame,开发调试时混着来——中间结果注册视图,SQL 做验证。
4.4 在 SQL 中使用变量
SQL 语句里不能直接写 Python 变量,需要用字符串格式化:
target_date = "2026-08-21"
min_amount = 0
# 方式一:f-string(注意引号嵌套)
result = spark.sql(f"""
SELECT * FROM clean_orders
WHERE order_date = '{target_date}' AND amount > {min_amount}
""")
# 方式二:format
result = spark.sql("""
SELECT * FROM clean_orders
WHERE order_date = '{}' AND amount > {}
""".format(target_date, min_amount))
安全提醒:如果变量来自用户输入(比如 Web 表单),直接拼接有 SQL 注入风险。Spark SQL 本身不支持参数化查询(PreparedStatement),需要在应用层做输入校验和转义。
五 常用内置函数速查
Spark SQL 的内置函数和 DataFrame API 的 pyspark.sql.functions 是同一套东西,只是调用方式不同。SQL 里写 COUNT(*),DataFrame 里写 count("*"),底层都是同一个 Catalyst 表达式。
5.1 按功能分类
| 聚合 | COUNT, SUM, AVG, MAX, MIN, COLLECT_LIST, COLLECT_SET | SELECT city, COLLECT_SET(category) FROM orders GROUP BY city |
| 字符串 | CONCAT, CONCAT_WS, SUBSTRING, LENGTH, UPPER, LOWER, TRIM, SPLIT | SELECT CONCAT_WS('-', city, district) FROM users |
| 日期时间 | CURRENT_DATE, CURRENT_TIMESTAMP, DATEDIFF, DATE_ADD, DATE_SUB, DATE_FORMAT, TO_DATE, UNIX_TIMESTAMP | SELECT DATEDIFF(CURRENT_DATE, order_date) FROM orders |
| 数学 | ROUND, CEIL, FLOOR, ABS, POW, SQRT, EXP, LOG | SELECT ROUND(amount / 100, 2) FROM orders |
| 条件 | CASE WHEN, COALESCE, IF, IFNULL, NULLIF | SELECT COALESCE(city, '未知') FROM users |
| 窗口 | ROW_NUMBER, RANK, DENSE_RANK, LAG, LEAD, SUM() OVER, AVG() OVER | SELECT ROW_NUMBER() OVER (PARTITION BY city ORDER BY amount DESC) FROM orders |
| JSON | GET_JSON_OBJECT, FROM_JSON, TO_JSON, JSON_TUPLE | SELECT GET_JSON_OBJECT(extra_info, '$.payment') FROM orders |
| 数组 | SIZE, ARRAY_CONTAINS, EXPLODE, ARRAY_JOIN | SELECT SIZE(category_list) FROM user_categories |
| Map | MAP_KEYS, MAP_VALUES, ELEMENT_AT | SELECT MAP_KEYS(category_amount) FROM user_stats |
5.2 与 DataFrame API 的命名差异
SQL 函数名通常是全大写,DataFrame API 是驼峰或下划线。大部分一一对应,少数有差异:
| COUNT(*) | count("*") 或 count(lit(1)) | * 在 DataFrame 里要用字符串 |
| COLLECT_LIST(x) | collect_list("x") | 名称一致 |
| GET_JSON_OBJECT(col, '$.key') | get_json_object("col", "$.key") | 名称一致 |
| ROW_NUMBER() | row_number() | 名称一致 |
| DATEDIFF(end, start) | datediff("end", "start") | 参数顺序一致 |
| DATE_FORMAT(date, 'yyyy-MM') | date_format("date", "yyyy-MM") | 名称一致 |
| CASE WHEN … THEN … END | when(…).otherwise(…) | 语法结构不同 |
| CAST(x AS DOUBLE) | col("x").cast("double") | 类型转换写法不同 |
5.3 复杂类型函数:Spark SQL 与 MySQL 的区别
看到 STRUCT 你可能会想:MySQL 里好像没有这个函数?
你的直觉是对的。Spark SQL 和传统关系型数据库(MySQL、PostgreSQL)不一样——它原生支持复杂数据类型。
传统 MySQL 每列只能存一个标量值(数字、字符串、日期)。但大数据场景下经常遇到半结构化数据(JSON、Parquet),里面嵌套着对象、数组、Map。Spark SQL 直接支持三类复杂类型:
| Struct | STRUCT(city, amount) | 把多个字段打包成一个对象 | ❌ 没有 |
| Array | ARRAY(1, 2, 3) | 数组/列表 | ❌ 没有 |
| Map | MAP('a', 1, 'b', 2) | 键值对 | ❌ 没有 |
为什么 Spark SQL 要支持这些?因为你经常会碰到这样的数据:
{
"user_id": 1001,
"address": {
"city": "上海",
"district": "浦东"
},
"tags": ["VIP", "新用户"],
"stats": {"login_days": 30, "order_count": 15}
}
MySQL 存这种数据得拆成三四张表,Spark SQL 直接读直接查。
下面分别看看这三类复杂类型在 SQL 里怎么用。
Struct(结构体)
— 构造 Struct
SELECT STRUCT(city, amount) as info FROM orders
— 访问 Struct 字段
SELECT info.city, info.amount FROM (
SELECT STRUCT(city, amount) as info FROM orders
)
# DataFrame 等价写法
from pyspark.sql.functions import struct
orders.select(struct("city", "amount").alias("info"))
# 访问字段
orders.select(struct("city", "amount").alias("info")).select("info.city", "info.amount")
Array(数组)
— 构造数组
SELECT ARRAY(category, brand) as tags FROM products
— 取数组第 N 个元素(注意从 1 开始)
SELECT tags[1] as first_tag FROM (
SELECT ARRAY(category, brand) as tags FROM products
)
— 数组长度
SELECT SIZE(tags) FROM products
— 数组合并(一行变多行)
SELECT EXPLODE(tags) as tag FROM products
Map(键值对)
— 构造 Map
SELECT MAP('city', city, 'tier', user_tier) as user_info FROM clean_orders
— 按 key 取值
SELECT user_info['city'] FROM (
SELECT MAP('city', city, 'tier', user_tier) as user_info FROM clean_orders
)
— 所有 key / 所有 value
SELECT MAP_KEYS(user_info), MAP_VALUES(user_info) FROM ...
六 UDF:当内置函数不够用的时候
6.1 为什么需要 UDF
内置函数覆盖了 90% 的场景,但总有例外。回到我们的日报任务,假设业务方加了两个需求:
第一个需求用内置函数也能写(CONCAT(SUBSTRING(phone,1,3), '****', SUBSTRING(phone,8,4))),但第二个需求(中文分词)内置函数搞不定——需要调用 jieba 这样的 Python 库。
这时候就需要 UDF(User Defined Function):用 Python 写自定义逻辑,然后在 SQL 或 DataFrame 中使用。
先理清楚层次:你看到的到底是哪一层的东西?
学 UDF 的时候很容易混淆,因为代码里同时出现了好几种不同层面的东西:
- Python 语言层:普通的 Python 函数、变量、if/else。比如 def mask_phone(phone): 这个函数,跟 Spark 一点关系都没有,你在任何 Python 脚本里都能写。
- DataFrame API 层:Spark 提供的 Python 接口。df.filter()、col()、udf() 这些都是。它把 Spark 的能力封装成了 Python 函数。
- SQL 语法层:spark.sql("SELECT …") 里写的 SQL 语句。
UDF 的作用就是搭桥:把 Python 层的函数,包装成 DataFrame 层和 SQL 层都能用的东西。
Python 函数(mask_phone)
│
├── 用 udf() 包装 ──▶ DataFrame API 能用
│
└── 用 register() 注册 ──▶ SQL 能用
6.2 定义和使用 UDF
from pyspark.sql.functions import udf, col
from pyspark.sql.types import StringType
# 1. 定义普通 Python 函数
def mask_phone(phone):
"""手机号脱敏:138****5678"""
if phone and len(phone) == 11:
return phone[:3] + "****" + phone[7:]
return phone
# 2. 注册成 UDF
mask_phone_udf = udf(mask_phone, StringType())
# 3. 在 DataFrame 中使用
users.withColumn("phone_masked", mask_phone_udf(col("phone")))
# 4. 注册成 SQL 函数,在 SQL 中使用
spark.udf.register("mask_phone", mask_phone, StringType())
spark.sql("SELECT mask_phone(phone) FROM users")
6.3 UDF 的性能问题
这是 UDF 最重要的知识点:Python UDF 很慢。
# 这段代码看起来没问题,但性能很差
users.withColumn("phone_masked", mask_phone_udf(col("phone")))
为什么内置函数快?因为它根本没离开 JVM
你可能会想:内置函数不也是函数吗,为啥它就快?
关键区别是:内置函数是用 Scala/Java 写的,编译在 JVM 里;数据从一开始就在 JVM 内存里,处理完还在 JVM 内存里。从头到尾没有跨进程,没有序列化,没有数据搬运。
┌──────────────────────────────────────────┐
│ JVM (Executor) │
│ │
│ 数据在 JVM 堆内存里 │
│ │ │
│ ▼ │
│ 内置函数(SUM / COUNT / WHEN…) │
│ │ │
│ ▼ │
│ 结果还在 JVM 堆内存里 │
│ │
└──────────────────────────────────────────┘
Python 进程:完全没参与
你在 Python 里写 col("amount") + 1,不是真的在 Python 里做加法。它只是构建了一个"表达式对象",告诉 Catalyst 优化器"我要做这个操作"。真正执行的时候,加法是 JVM 里的代码在 JVM 内存里直接跑的。
Python 在这里只是个"指挥官"——它告诉 JVM 要做什么,真正干活的是 JVM。内置函数就是 JVM 里的"工人",而 Python UDF 是 Python 里的"外援工人"——数据得搬过去让外援做,做完再搬回来。
更进一步说,不光内置函数是这样,整个 DataFrame API 和 SQL 都是这样:你写的代码只是在描述"我要什么结果",不是在指挥"一步步怎么干"。Spark 会把你的描述翻译成一棵"逻辑计划树",然后交给 Catalyst 优化器去想办法最快地得到结果。这就是 Spark 的声明式编程 + 懒执行——.filter()、.select() 这些操作只是搭积木,等到 .show() 或 .write() 时才真正动工。
为什么 Python UDF 慢?因为数据要来回搬家
Python UDF 的函数体在 Python 进程里。数据本来在 JVM 里,要处理必须跨进程搬运:
┌──────────────────┐ 序列化 ┌──────────────────┐
│ JVM (Executor) │ ──────────────▶ │ Python 进程 │
│ │ 每一行数据 │ │
│ 数据在 JVM 里 │ │ 逐行调用 UDF │
│ │ ◀────────────── │ │
│ │ 序列化结果 │ 计算完返回 │
└──────────────────┘ └──────────────────┘
慢的三个原因:
官方基准测试:同样的逻辑,Python UDF 比内置函数慢 10-100 倍。

6.4 Pandas UDF(向量化 UDF)
Spark 2.3+ 引入了 Pandas UDF,利用 Apache Arrow 在 JVM 和 Python 之间批量传输数据,用 Pandas 的向量化操作代替逐行处理。
from pyspark.sql.functions import pandas_udf, col
from pyspark.sql.types import StringType
import pandas as pd
# Pandas UDF:输入输出都是 Pandas Series
@pandas_udf(StringType())
def mask_phone_pandas(phone_series: pd.Series) –> pd.Series:
return phone_series.str[:3] + "****" + phone_series.str[7:]
# 使用方式和普通 UDF 一样
users.withColumn("phone_masked", mask_phone_pandas(col("phone")))
"向量化"和"批量"不是一回事
你可能会想:不就是批量传数据吗,为啥叫"向量化"?
其实 Pandas UDF 的加速来自两层优化的叠加:
普通 Python UDF:
JVM ──行1──▶ Python(逐行解释执行)──结果1──▶ JVM
JVM ──行2──▶ Python(逐行解释执行)──结果2──▶ JVM
…
每行都要:序列化 + 跨进程 + Python 解释器 + 反序列化
Pandas UDF:
JVM ──一批数据(Arrow)──▶ Python(C 层数组操作)──一批结果──▶ JVM
一批数据只走一次跨进程,计算用 C 向量化,快很多
向量化的本质是:把 Python 层的 for 循环,换成了底层 C 层的数组操作。 Python 的 for 循环很慢(每行都要经过解释器、类型检查、函数调用),而 Pandas / NumPy 的操作底层是 C 写的,直接操作连续内存。
性能对比:
| Python UDF | 逐行序列化 | 逐行 Python 解释执行 | 1x(基准) |
| Pandas UDF | Arrow 批量 | C 层向量化计算 | 10-50x |
| 内置函数 | 无跨进程 | JVM 原生 | 100x |
原则:能用内置函数就别用 UDF;必须用 UDF 时优先 Pandas UDF;Python UDF 只在小数据量或无法避免时使用。
- 内置函数:Python 这边只是"发号施令"——用 Python 函数构造一个表达式描述,发给 JVM。JVM 认识这个表达式(因为是 Spark 自己预置的),直接在 JVM 里执行,数据不搬家。
- UDF:JVM 不认识你写的 Python 函数(里面可能调 jieba、可能写 if-else、可能有任何 Python 语法),JVM 没法执行。所以只能把数据序列化发给 Python 进程,让 Python 解释器自己执行你的函数,执行完再把结果序列化传回来。
6.5 什么时候不得不用 UDF
尽管性能差,以下场景还是得用:
既然 Spark 是 JVM 的,为啥不用 Java/Scala 写?
这是个好问题。理论上用 Scala/Java 写确实"更原生"、性能更好。但现实中大多数人用 PySpark,原因很实际:
- Python 生态碾压:数据科学、机器学习、NLP、可视化……Python 库的丰富程度是 JVM 比不了的。中文分词用 jieba,训练模型用 scikit-learn / PyTorch,这些都是 Python 里一行 import 的事,JVM 里要么没有,要么难用得多。
- 开发效率高:同样的逻辑,Python 代码量通常是 Java 的 1/2 到 1/3。不用编译、语法简洁、写了就能跑,探索式分析的迭代速度快很多。
- 团队技能结构:大部分数据团队(分析师、算法工程师、数据科学家)的背景是 Python,不是 Java/Scala。
- 不是所有场景都要极致性能:很多 UDF 用在一天跑一次的离线任务或探索式分析里,慢个几分钟完全能接受,但开发省了几小时。
PySpark 的核心价值就是:用 Python 的便利性,享受 Spark 的分布式计算能力。 两边的好处都占了。当然,如果是性能要求极高的核心链路,或者要写 Spark 底层扩展,还是得用 Scala/Java。
回到我们的商品标题分词需求:
# 示例:在 UDF 中调用外部库做文本分词
import jieba
from pyspark.sql.functions import pandas_udf, col
from pyspark.sql.types import ArrayType, StringType
import pandas as pd
@pandas_udf(ArrayType(StringType()))
def chinese_tokenize(text_series: pd.Series) –> pd.Series:
return text_series.apply(lambda x: list(jieba.cut(x)) if x else [])
# 商品标题分词
products.withColumn("title_tokens", chinese_tokenize(col("title")))
这种涉及外部 Python 库的需求,内置函数搞不定,就只能用 UDF。好在 Pandas UDF 的性能已经比普通 Python UDF 好很多了。
6.6 UDF 在 SQL 中的注册与使用
# 注册
spark.udf.register("mask_phone", mask_phone, StringType())
# 在 SQL 中使用
spark.sql("""
SELECT user_id, mask_phone(phone) as phone_masked
FROM users
WHERE city = '上海'
""")
# 也可以注册 Pandas UDF
spark.udf.register("mask_phone_pandas", mask_phone_pandas)
注意:spark.udf.register() 注册的函数只在当前 SparkSession 有效。如果通过 JDBC/Thrift Server 连接,需要在服务器端注册。
七 进阶案例:帕累托分析与月度环比(SQL 版)
用 SQL 重写第三篇的两个经典分析,体会 SQL 写法的优势。数据还是基于我们日报任务里的 clean_orders 视图。
7.1 帕累托分析(SQL 版)
# 先按用户聚合消费金额
user_spending = spark.sql("""
SELECT user_id, SUM(amount) as total_amount
FROM clean_orders
GROUP BY user_id
""")
user_spending.createOrReplaceTempView("user_spending")
# SQL 版帕累托分析
pareto_sql = spark.sql("""
WITH ranked AS (
SELECT
user_id,
total_amount,
— 全局累计求和:按金额降序,从第一行到当前行
SUM(total_amount) OVER (ORDER BY total_amount DESC) as cumsum_amount,
— 全局总计
SUM(total_amount) OVER () as grand_total
FROM user_spending
)
SELECT
user_id,
total_amount,
cumsum_amount,
ROUND(cumsum_amount / grand_total, 4) as cumsum_ratio
FROM ranked
WHERE cumsum_amount / grand_total <= 0.8
""")
# 验证结果
pareto_sql.show(5)
与 DataFrame 版对比:
# DataFrame 版(第三篇)
from pyspark.sql.window import Window
from pyspark.sql.functions import spark_sum, round as spark_round
window_cumsum = Window.orderBy(col("total_amount").desc())
window_total = Window.partitionBy()
pareto_df = (
user_spending
.withColumn("cumsum_amount", spark_sum("total_amount").over(window_cumsum))
.withColumn("grand_total", spark_sum("total_amount").over(window_total))
.withColumn("cumsum_ratio", spark_round(col("cumsum_amount") / col("grand_total"), 4))
.filter(col("cumsum_amount") / col("grand_total") <= 0.8)
)
SQL 版的优势:WITH 子句(CTE)让逻辑分层更清晰,窗口函数的写法也更贴近标准 SQL 语法。
DataFrame 版的优势:类型安全、可链式调用、便于单元测试。
7.2 月度环比(SQL 版)
monthly_sql = spark.sql("""
WITH monthly AS (
SELECT
city,
YEAR(order_time) as yr,
MONTH(order_time) as mo,
SUM(amount) as monthly_amount
FROM clean_orders
GROUP BY city, YEAR(order_time), MONTH(order_time)
)
SELECT
city,
yr,
mo,
monthly_amount,
LAG(monthly_amount, 1) OVER (PARTITION BY city ORDER BY yr, mo) as prev_month_amount,
ROUND(
(monthly_amount – LAG(monthly_amount, 1) OVER (PARTITION BY city ORDER BY yr, mo))
/ LAG(monthly_amount, 1) OVER (PARTITION BY city ORDER BY yr, mo),
4
) as mom_growth_rate
FROM monthly
ORDER BY city, yr, mo
""")
monthly_sql.show()
SQL 的重复问题:LAG(…) 写了三遍,SQL 没有变量复用机制。DataFrame 可以先把 lag 结果存成列再引用,SQL 只能重复写或嵌套子查询。
7.3 最佳实践模式:DataFrame 做管道,SQL 做报表
实际工作中的典型模式,也是本篇贯穿始终的思路:
# 开发用 DataFrame 写每日 ETL 管道
def daily_etl(target_date):
# 复杂的数据清洗、关联、派生用 DataFrame
clean_data = (
spark.table("dw.orders")
.filter(col("order_date") == target_date)
.filter(col("order_status") == "completed")
.join(spark.table("dw.users"), "user_id")
.join(spark.table("dw.products"), "product_id")
.withColumn("user_tier", when(col("amount") > 1000, "VIP").otherwise("普通"))
)
# 注册成视图,供分析师和 BI 工具查询
clean_data.createOrReplaceTempView("daily_orders")
return clean_data
# 分析师用 SQL 做日报表
report = spark.sql("""
SELECT
city,
user_tier,
COUNT(DISTINCT user_id) as active_users,
SUM(amount) as gmv,
AVG(amount) as avg_order_value
FROM daily_orders
GROUP BY city, user_tier
ORDER BY gmv DESC
""")
# BI 工具通过 JDBC 直接查视图
# Tableau / FineBI -> jdbc:spark://host:10000 -> SELECT * FROM daily_orders
八 小结
| Spark SQL 基础 | spark.sql() 直接查 Paimon 表和文件,返回 DataFrame |
| 临时视图 | createOrReplaceTempView() 注册,作用域为当前 Session,是 DataFrame → SQL 的桥梁 |
| 全局临时视图 | createOrReplaceGlobalTempView() 注册,跨 Session 共享,查询加 global_temp. 前缀 |
| SQL vs DataFrame | 底层同一优化器,性能一致;SQL 适合报表和临时查询,DataFrame 适合 ETL 管道 |
| 内置函数 | 覆盖 90% 场景,SQL 和 DataFrame API 是同一套函数的不同写法 |
| UDF | 内置函数不够用时才用;Python UDF 慢(逐行 + 跨进程),Pandas UDF 快(向量化 + Arrow) |
| 最佳实践 | 开发用 DataFrame 写管道,中间结果注册视图;分析师用 SQL 查询和报表 |
下一篇将进入性能调优:Shuffle 机制、缓存策略、数据倾斜处理,让任务跑得更快更稳。




