京公网安备 11010802034615号
经营许可证编号:京B2-20210330

做数据聚合时,PySpark的groupBy()确实能完成统计,这也是它的本职工作。但它有一个根本性局限:每一组数据,最终只能返回一行结果。不管你对上千、上百万行数据做SUM求和,输出都只有一行。
很多场景下这个结果够用,但有时我们希望保留原始明细,同时附带分组聚合指标。这时,就该轮到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|
+--------------+----------------+----------+--------+------+
每一行代表一笔交易。接下来我们使用窗口函数,保留全部交易行,同时做数据分析。
窗口定义了:计算当前行指标时,需要参考哪些数据行。 窗口定义一般包含三部分:
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行,每家门店一行汇总。
✅ 一句话分清两者:
排名是窗口函数最常用场景。比如:按门店,把订单金额从高到低排名。 先定义窗口:按门店分区,金额降序排列
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,依次顺延。
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。
”
排名最经典用法:取出每个分组的头部数据。例如,取出每家门店金额最高的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)拿的是全数据集前两名,不是每家门店前两名。
这类写法非常适合业务问题:
累计求和:把当前行+前面所有行累加。我们计算每家门店随时间的累计销售额。 窗口规则:
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():按行下标选取,.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个优化习惯:
recent_sales = sales.filter(
F.col("sale_date") >= F.lit("2026-01-05"))
window_input = sales.select(
"transaction_id",
"store",
"sale_date",
"amount",
)
Window.orderBy(F.col("amount").desc()) 是全局排序,不是分组内排序,和按门店分区排名完全不一样。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排序,大数据量注意看执行计划,优先前置过滤
弄懂分区、排序、窗口帧三者的配合,窗口函数就能轻松解决绝大多数数据工程常见统计难题。

做数据聚合时,PySpark的groupBy()确实能完成统计,这也是它的本职工作。但它有一个根本性局限:每一组数据,最终只能返回一行 ...
2026-10-01热力地图是数据可视化中极具辨识度与实用性的空间分析图表,结合地理空间维度与数据密度特征,通过颜色深浅、色阶渐变直观展示数 ...
2026-09-30 很多数据分析师做过按月份的销售额趋势图,画过按天的流量折线图,但当被问到“时间序列和普通数据有什么本质区别”“季节性 ...
2026-09-30同样是“银行数据岗”,在国有大行总行数据中心、在一家城商行的零售部、在银行系金融科技子公司、在保险公司,工作内容、成长节 ...
2026-09-29在数据分析与统计学研究中,数据往往不是独立存在的,不同变量之间普遍存在相互关联、相互影响的关系。相关性统计分析是挖掘变量 ...
2026-09-29 导读:大多数人只把 dataclasses 当成偷懒工具,用来少写 __init__、__repr__ 这类魔法方法。但它的能力远不止于此。本文带 ...
2026-09-29 很多数据分析师能熟练地计算指标、搭建标签体系,但当被问到“画像到底在解决什么问题”“画像和标签是什么关系”“画像如何 ...
2026-09-29在MySQL数据库运维与业务开发中,行业普遍存在“数据达到千万级就必须分表”的说法。但在实际生产环境中,千万条数据并不是强制 ...
2026-09-28CDA数据分析师 出品 作者:李诗怡 1. 5W1H 分析法 定义:经典系统性思维框架,通过六个核心维度对问题进行全方位拆解与剖析,确 ...
2026-09-28 很多分析师在设计标签时思路清晰,但真到落地环节却面临“数据在手,不知如何转化为可用标签”的困境:或因加工方式选择不当 ...
2026-09-28CDA数据分析师 出品 作者:李诗怡 1. 用户标签体系 定义: 通过一系列高度精炼的特征标识,对用户属性、行为与偏好进行量化刻画 ...
2026-09-24Pandas是Python生态中用于表格数据处理的核心库,广泛应用于数据清洗、统计运算、报表输出、数据分析建模等场景。在处理极大数值 ...
2026-09-24随着数字经济快速发展,数据已成为核心生产要素,各行各业的业务沉淀、用户行为、设备运行、市场交易均产生海量数据。数据处理作 ...
2026-09-24 很多分析师每天和数据打交道,但当被问到“标签是什么”“标签和指标有什么区别”“标签体系如何设计”时,却常常答不上来。 ...
2026-09-24在时序数据分析中,大部分业务数据并非持续平稳变化,而是会在某些时间节点出现突然抬升、断崖下跌、趋势反转、波动异变等现象, ...
2026-09-23在统计学与数据分析中,研究多组数据差异最常用的方法为单因素方差分析与事后多重比较。很多数据分析初学者容易混淆两者功能,认 ...
2026-09-23 很多数据分析师每天都在写 SQL,但当被问到“DQL 的本质是什么”“SELECT 子句的书写顺序与执行顺序为何不同”“INNER JOIN ...
2026-09-23 很多数据分析师写过无数个SELECT查询,但当被问到“如何新建一张表来固化中间数据”“创建视图和创建物理表有什么区别”“视 ...
2026-09-22CDA数据分析师 出品 作者:李诗怡 1. 金字塔原理 定义: 一种“先总后分、先结论后原因”的思考和表达方式。顶层为核心观点,中 ...
2026-09-22数据收集是数据分析、数据挖掘与数字化运营的源头工作,数据收集的完整性、准确性、时效性直接决定后续数据分析结果的可信度与业 ...
2026-09-21