热线电话:13121318867

登录
首页大数据时代PySpark窗口函数实战入门:为什么groupBy不够用?
PySpark窗口函数实战入门:为什么groupBy不够用?
2026-10-01
收藏

做数据聚合时,PySpark的groupBy()确实能完成统计,这也是它的本职工作。但它有一个根本性局限:每一组数据,最终只能返回一行结果。不管你对上千、上百万行数据做SUM求和,输出都只有一行。

很多场景下这个结果够用,但有时我们希望保留原始明细,同时附带分组聚合指标。这时,就该轮到PySpark窗口函数登场了。

窗口函数可以基于一组关联记录做计算,且不会把多条记录压缩成单行。也就是说,你既能保留每一条交易明细,又能拿到所属分组的汇总统计。 适合实现分组内排名、累计求和、和上一条记录对比、组内占比计算等需求。

案例基于销售数据集演示,这套写法同样适用于事件日志、财务记录、用户行为、传感器时序等各类分组有序数据。

”

目录

  1. PySpark环境准备
  2. 构建示例数据集
  3. 什么是窗口(Window)
  4. 分组内行排名
  5. row_number、rank、dense_rank三者区别
  6. 提取每组TopN数据
  7. 计算累计总和
  8. 当前行与上一行对比
  9. 计算单行在分组总和中的占比
  10. 移动平均值计算
  11. rowsBetween 与 rangeBetween 窗口帧
  12. 复用窗口定义
  13. 性能优化要点
  14. 高频踩坑误区
  15. 综合案例:多指标联合计算
  16. 总结

PySpark环境准备

如果本地还未安装PySpark,新建项目文件夹,用uv安装依赖:

mkdir pyspark-windows
cd pyspark-windows
uv init
uv venv
.venvScriptsactivate
uv pip install pyspark

创建Spark会话,验证环境是否正常:

import os
import sys
os.environ["PYSPARK_PYTHON"] = sys.executable
os.environ["PYSPARK_DRIVER_PYTHON"] = sys.executable
from pyspark.sql import SparkSession

spark = (
    SparkSession.builder
    .master("local[*]")
    .appName("sales-analysis")
    .config("spark.pyspark.python", sys.executable)
    .config("spark.pyspark.driver.python", sys.executable)
    .getOrCreate()
)

local[*]代表本地模式运行,自动调用本机所有CPU核心,跑案例无需搭建Spark集群。 保存代码后,执行命令运行脚本:

spark-submit spark_example.py

构建示例数据集

数据集包含3家门店多日销售记录:

from datetime import date
from pyspark.sql import functions as F
from pyspark.sql import types as T
from pyspark.sql.window import Window

sales_data = [
    (1, "London_Store", date(2026, 1, 2), "Laptop", 1200.00),
    (2, "London_Store", date(2026, 1, 3), "Monitor", 350.00),
    (3, "London_Store", date(2026, 1, 5), "Keyboard", 90.00),
    (4, "London_Store", date(2026, 1, 8), "Laptop", 1350.00),
    (5, "Manchester_Store", date(2026, 1, 2), "Monitor", 320.00),
    (6, "Manchester_Store", date(2026, 1, 4), "Laptop", 1100.00),
    (7, "Manchester_Store", date(2026, 1, 6), "Mouse", 45.00),
    (8, "Manchester_Store", date(2026, 1, 9), "Laptop", 1250.00),
    (9, "Bristol_Store", date(2026, 1, 3), "Keyboard", 85.00),
    (10, "Bristol_Store", date(2026, 1, 4), "Monitor", 300.00),
    (11, "Bristol_Store", date(2026, 1, 7), "Laptop", 1050.00),
    (12, "Bristol_Store", date(2026, 1, 10), "Monitor", 330.00),
]
sales_schema = T.StructType(
    [
        T.StructField("transaction_id", T.IntegerType(), False),
        T.StructField("store", T.StringType(), False),
        T.StructField("sale_date", T.DateType(), False),
        T.StructField("product", T.StringType(), False),
        T.StructField("amount", T.DoubleType(), False),
    ]
)
sales = spark.createDataFrame(sales_data, schema=sales_schema)
sales.orderBy("store", "sale_date").show()

输出结果:

+--------------+----------------+----------+--------+------+
|transaction_id|           store| sale_date| product|amount|
+--------------+----------------+----------+--------+------+
|             9|   Bristol_Store|2026-01-03|Keyboard|  85.0|
|            10|   Bristol_Store|2026-01-04| Monitor| 300.0|
|            11|   Bristol_Store|2026-01-07|  Laptop|1050.0|
|            12|   Bristol_Store|2026-01-10| Monitor| 330.0|
|             1|    London_Store|2026-01-02|  Laptop|1200.0|
|             2|    London_Store|2026-01-03| Monitor| 350.0|
|             3|    London_Store|2026-01-05|Keyboard|  90.0|
|             4|    London_Store|2026-01-08|  Laptop|1350.0|
|             5|Manchester_Store|2026-01-02| Monitor| 320.0|
|             6|Manchester_Store|2026-01-04|  Laptop|1100.0|
|             7|Manchester_Store|2026-01-06|   Mouse|  45.0|
|             8|Manchester_Store|2026-01-09|  Laptop|1250.0|
+--------------+----------------+----------+--------+------+

每一行代表一笔交易。接下来我们使用窗口函数,保留全部交易行,同时做数据分析。

什么是窗口Window

窗口定义了:计算当前行指标时,需要参考哪些数据行。 窗口定义一般包含三部分:

  • partitionBy():对数据分组
  • orderBy():组内行排序
  • rowsBetween() / rangeBetween():相对当前行划定计算区间(窗口帧)

基础窗口示例:

store_window = Window.partitionBy("store")

以门店分组,伦敦门店的交易在一个分区,曼彻斯特、布里斯托各自独立分区。

搭配聚合函数使用:

sales_with_store_total = sales.withColumn(
    "store_total",
    F.sum("amount").over(store_window),
)
sales_with_store_total.orderBy("store", "sale_date").show()

输出里每一行交易,都会附带所属门店的总销售额,原始明细一条都不会丢失。

对比groupBy写法:

sales.groupBy("store").agg(
    F.sum("amount").alias("store_total")
).show()

groupBy执行后,只返回3行,每家门店一行汇总。

✅ 一句话分清两者:

  • groupBy:每组只输出汇总行,明细丢失
  • 窗口函数:保留全部原始明细,附加分组统计

分组内行排名

排名是窗口函数最常用场景。比如:按门店,把订单金额从高到低排名。 先定义窗口:按门店分区,金额降序排列

sales_rank_window = (Window.partitionBy("store").orderBy(F.col("amount").desc()))

使用row_number编号:

ranked_sales = sales.withColumn(
    "sale_rank",
    F.row_number().over(sales_rank_window),
)
ranked_sales.orderBy("store", "sale_rank").show()

同门店内,金额最高订单rank=1,依次顺延。

row_number、rank、dense_rank 三者区别

ranking_comparison = (
    sales
    .withColumn(
        "row_number",
        F.row_number().over(sales_rank_window),
    )
    .withColumn(
        "rank",
        F.rank().over(sales_rank_window),
    )
    .withColumn(
        "dense_rank",
        F.dense_rank().over(sales_rank_window),
    )
)
ranking_comparison.show()

核心差异(存在并列值时才会体现)

  • row_number():永远生成唯一连续序号,并列数据也强行区分编号
  • rank():并列值给相同排名,后面序号留空位
  • dense_rank():并列值给相同排名,后面序号不留空位

举例:金额100、100、80

amount row_number rank dense_rank
100 1 1 1
100 2 1 1
80 3 3 2

选择建议:需要严格唯一序号用row_number;并列同名次用rank/dense_rank。

”

提取每组TopN记录

排名最经典用法:取出每个分组的头部数据。例如,取出每家门店金额最高的2笔订单:

top_two_sales_per_store = (
    sales
    .withColumn(
        "sale_rank",
        F.row_number().over(sales_rank_window),
    )
    .filter(F.col("sale_rank") <= 2)
    .orderBy("store", "sale_rank")
)
top_two_sales_per_store.show()

⚠️ 不要直接全局limit!sales.orderBy(F.col("amount").desc()).limit(2)拿的是全数据集前两名,不是每家门店前两名。

这类写法非常适合业务问题:

  • 每个客户金额最高的5条订单
  • 每个区域销量最好的3款商品
  • 每台设备最近两条事件记录

计算累计总和

累计求和:把当前行+前面所有行累加。我们计算每家门店随时间的累计销售额。 窗口规则:

  1. 按门店分区
  2. 按交易日期、交易ID排序
  3. 区间:分区第一行 ~ 当前行
running_total_window = (
    Window
    .partitionBy("store")
    .orderBy("sale_date", "transaction_id")
    .rowsBetween(
        Window.unboundedPreceding,
        Window.currentRow,
    )
)

sales_with_running_total = sales.withColumn(
    "running_store_total",
    F.sum("amount").over(running_total_window),
)
sales_with_running_total.orderBy("store", "sale_date").show()

每家门店第一笔交易作为累计起点,后续每一笔持续累加。

排序增加transaction_id:同一天多条交易时,用来固定顺序,避免结果不稳定。

”

当前行与上一行对比

lag()函数,提取窗口内上一行数据,非常适合时序指标环比。

store_date_window = (
    Window
    .partitionBy("store")
    .orderBy("sale_date", "transaction_id")
)

sales_with_previous_amount = (
    sales
    .withColumn(
        "previous_amount",
        F.lag("amount").over(store_date_window),
    )
    .withColumn(
        "change_from_previous",
        F.col("amount") - F.col("previous_amount"),
    )
)
sales_with_previous_amount.orderBy("store", "sale_date").show()

门店第一条交易没有上一条记录,字段返回NULL。

还可以对比日期,计算两次交易间隔天数:

sales_with_previous_date = (
    sales
    .withColumn(
        "previous_sale_date",
        F.lag("sale_date").over(store_date_window),
    )
    .withColumn(
        "days_since_previous_sale",
        F.datediff("sale_date", "previous_sale_date"),
    )
)
sales_with_previous_date.show()

可以用来识别用户活跃断层、事件延迟、指标波动。

配套lead(),取下一行数据:

sales_with_next_date = sales.withColumn(
    "next_sale_date",
    F.lead("sale_date").over(store_date_window),
)
sales_with_next_date.show()

计算单行在分组总和中的占比

窗口聚合可以轻松算出单条记录占本组总和百分比。 本例:每一笔销售额,占所属门店总营收比例。

store_total_window = Window.partitionBy("store")
sales_with_share = (
    sales
    .withColumn(
        "store_total",
        F.sum("amount").over(store_total_window),
    )
    .withColumn(
        "share_of_store_sales",
        F.col("amount") / F.col("store_total"),
    )
)
sales_with_share.select(
    "store",
    "transaction_id",
    "amount",
    "store_total",
    F.round("share_of_store_sales", 3).alias(
        "share_of_store_sales"
    ),
).orderBy("store", F.col("amount").desc()).show()

不用单独聚合再join回原表,一步完成。 同类场景:员工薪资占部门人力成本、单品销售额占品类总额、单笔消费占客户总消费。

移动平均值计算

累计求和会包含分区全部前置数据;移动窗口只取附近有限行。 示例窗口:当前交易 + 前面2笔交易,共3条记录计算移动均值:

moving_average_window = (
    Window
    .partitionBy("store")
    .orderBy("sale_date", "transaction_id")
    .rowsBetween(-2, Window.currentRow)
)

sales_with_moving_average = sales.withColumn(
    "three_sale_average",
    F.avg("amount").over(moving_average_window),
)
sales_with_moving_average.orderBy("store", "sale_date").show()

注意:这是按行数窗口,不是时间窗口。哪怕门店某天密集多单,也只会取最近3行,不是最近3天。

”

rowsBetween 和 rangeBetween 窗口帧

窗口帧用来划定参与计算的行范围。

  • rowsBetween():按行下标选取,.rowsBetween(-3, Window.currentRow)代表当前行+往前3行,一共4行
  • rangeBetween():基于排序字段的数值范围,而不是行数。

时间范围窗口示例(7天窗口):

seconds_in_seven_days = 7 * 24 * 60 * 60
seven_day_window = (
    Window
    .partitionBy("store")
    .orderBy(F.col("sale_timestamp").cast("long"))
    .rangeBetween(-seconds_in_seven_days, 0)
)

???? 选择口诀:固定N条记录用rowsBetween;固定时间/数值区间用rangeBetween。 时间窗口要格外小心,排序字段需要转为统一单位的数字。

复用窗口定义

窗口定义本身不会修改DataFrame,它只是描述分组、排序、区间规则。一次定义,多处复用,代码可读性大幅提升:

store_total_window = Window.partitionBy("store")
store_date_window = (
    Window
    .partitionBy("store")
    .orderBy("sale_date", "transaction_id")
)
running_total_window = (
    store_date_window
    .rowsBetween(
        Window.unboundedPreceding,
        Window.currentRow,
    )
)

analysed_sales = (
    sales
    .withColumn(
        "store_total",
        F.sum("amount").over(store_total_window),
    )
    .withColumn(
        "previous_amount",
        F.lag("amount").over(store_date_window),
    )
    .withColumn(
        "running_total",
        F.sum("amount").over(running_total_window),
    )
)
analysed_sales.show()

给窗口起清晰命名,后续排查逻辑更简单。

性能优化要点

窗口函数不是零成本运算。Spark需要移动、排序数据,把相同分区key的数据放到一起,这个shuffle操作在大数据量下开销很大。

查看执行计划,观察shuffle和sort:

analysed_sales.explain("formatted")

看到Exchange、Sort算子,就是发生了数据重排。

✅ 4个优化习惯:

  1. 提前过滤:窗口计算前先筛掉不需要的数据,减少shuffle数据量
recent_sales = sales.filter(
    F.col("sale_date") >= F.lit("2026-01-05"))
  1. 只选需要字段:剔除无关列再送入窗口计算
window_input = sales.select(
    "transaction_id",
    "store",
    "sale_date",
    "amount",
)
  1. 警惕数据倾斜:某个分区key行数远大于其他,会造成单个任务压力爆炸。选择分区key时,提前评估数据分布。
  2. 谨慎缓存:同一个窗口计算结果会被多次复用,才考虑.cache();不要无脑缓存所有中间表。

高频踩坑误区

  1. 忘记partitionBy:Window.orderBy(F.col("amount").desc()) 是全局排序,不是分组内排序,和按门店分区排名完全不一样。
  2. 排序字段不完整:多条记录排序字段值相同时,一定要增加辅助字段(例如transaction_id)固定顺序,否则结果不稳定。
  3. 误以为窗口函数会减少行数:窗口是新增字段,不会删减行。想要只保留TopN,必须在窗口后加filter过滤。
  4. 把行窗口当成时间窗口:rowsBetween(-6,0)代表7条记录,不是7天。按时间统计一定要用rangeBetween。

综合案例:多指标一次性计算

一次性在交易表追加门店总额、占比、累计销售额、环比变化、订单排名:

store_total_window = Window.partitionBy("store")
store_date_window = (
    Window
    .partitionBy("store")
    .orderBy("sale_date", "transaction_id")
)
running_total_window = (
    store_date_window
    .rowsBetween(
        Window.unboundedPreceding,
        Window.currentRow,
    )
)
store_rank_window = (
    Window
    .partitionBy("store")
    .orderBy(F.col("amount").desc())
)

sales_analysis = (
    sales
    .withColumn(
        "store_total",
        F.sum("amount").over(store_total_window),
    )
    .withColumn(
        "share_of_store_total",
        F.col("amount") / F.col("store_total"),
    )
    .withColumn(
        "running_store_total",
        F.sum("amount").over(running_total_window),
    )
    .withColumn(
        "previous_amount",
        F.lag("amount").over(store_date_window),
    )
    .withColumn(
        "change_from_previous",
        F.col("amount") - F.col("previous_amount"),
    )
    .withColumn(
        "sale_rank",
        F.row_number().over(store_rank_window),
    )
)
sales_analysis.orderBy("store", "sale_date").show()

原始交易明细全部保留,同时附带多维度分组统计指标。

总结

PySpark窗口函数,可以在保留原始明细行的前提下,跨关联行完成统计。当groupBy会丢失明细时,窗口函数就是最优解。

核心要点回顾:
✅ partitionBy划分独立分组
✅ 依赖顺序的计算,一定要搭配orderBy
✅ 通过窗口帧控制参与计算的行范围
✅ 排名函数用于组内对比;lag/lead用于时序前后对比
✅ sum/avg等聚合搭配窗口,实现分组汇总、占比、滑动统计
⚠️ 窗口会触发shuffle排序,大数据量注意看执行计划,优先前置过滤

弄懂分区、排序、窗口帧三者的配合,窗口函数就能轻松解决绝大多数数据工程常见统计难题。

推荐学习书籍 《CDA一级教材》适合CDA一级考生备考,也适合业务及数据分析岗位的从业者提升自我。完整电子版已上线CDA网校,累计已有10万+在读~ !

免费加入阅读:https://edu.cda.cn/goods/show/3151?targetId=5147&preview=0

数据分析师资讯
更多

OK
客服在线
立即咨询
客服在线
立即咨询