欢迎光临
我们一直在努力

【PySpark 学习笔记 五】SQL 协作:Spark SQL、临时视图、UDF

前三篇用 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 执行引擎),但对外开了两扇门:

DataFrame API 界面SQL 界面
入口 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 本篇的主线任务

本篇围绕一个具体任务展开:做一张城市维度的销售日报。

需要输出的报表长这样:

cityuser_tieractive_usersgmvavg_order_value
上海 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()

看起来还行?那再加几个需求试试:

  • 每个城市的 GMV 排名(窗口函数)
  • 对比昨日的环比增长(需要 self-join 或 LAG)
  • VIP 用户占比(需要先算总数,再算 VIP 数,再相除)
  • 加了窗口函数和多层聚合之后,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 起个别名 clean_orders
  • 把这个别名注册到当前 SparkSession 的 Catalog 里(同名则覆盖)
  • 关键点:临时视图不保存数据,它只是一个指向 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,反之亦然。

    操作DataFrame APISQL
    选列 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 是驼峰或下划线。大部分一一对应,少数有差异:

    SQLDataFrame 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 直接支持三类复杂类型:

    复杂类型Spark SQL 写法含义MySQL 里有吗
    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% 的场景,但总有例外。回到我们的日报任务,假设业务方加了两个需求:

  • 用户手机号脱敏展示(138****5678)
  • 商品标题分词,统计热门关键词
  • 第一个需求用内置函数也能写(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 │
    │ │ ◀────────────── │ │
    │ │ 序列化结果 │ 计算完返回 │
    └──────────────────┘ └──────────────────┘

    慢的三个原因:

  • 跨进程传输:数据要从 JVM Executor 序列化后传给 Python 进程,计算完再传回 JVM
  • 逐行处理:Python UDF 是逐行执行的,没有向量化优化
  • 无法下推:Catalyst 优化器看不懂 UDF 内部的逻辑,不能把它下推到数据源(比如不能利用 Paimon 的谓词下推)
  • 官方基准测试:同样的逻辑,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 的加速来自两层优化的叠加:

  • 批量传输(减少跨进程开销):不再逐行传数据,而是攒一批(比如几千到几万行)用 Arrow 格式一次性传过去,减少序列化/反序列化的次数
  • 向量化计算(减少 Python 循环开销):传过来的不是一行一行的值,而是一整列(Pandas Series)。phone_series.str[:3] 看着像 Python 代码,实际上底层是 C 实现的数组操作,直接操作连续内存,比 Python 逐行循环快得多
  • 普通 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

    尽管性能差,以下场景还是得用:

  • 调用外部 Python 库(没有 JVM 等价物)
  • 复杂业务逻辑,用 SQL 表达太冗长(比如 50 行 if-else 的风控规则)
  • 与外部系统交互(在 UDF 里调 HTTP API——不推荐,但有时不得不做)
  • 既然 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 机制、缓存策略、数据倾斜处理,让任务跑得更快更稳。

    赞(0)
    未经允许不得转载:171主机测评 » 【PySpark 学习笔记 五】SQL 协作:Spark SQL、临时视图、UDF
    分享到: 更多 (0)

    评论 抢沙发

    • 昵称 (必填)
    • 邮箱 (必填)
    • 网址