PySpark同公司支出按序匹配优化 替代collect低效实现方案
问题场景
你有两个支持多公司数据的PySpark DataFrame,结构如下:
df1 结构
| Company Name | UnderspendAmount |
|---|---|
| Pepsi | 2700.00 |
| Pepsi | 1200.00 |
df2 结构
| Company Name | OverspendAmount |
|---|---|
| Pepsi | 3600.00 |
| Pepsi | 1400.00 |
需求是针对同一个Company Name,将最高的UnderspendAmount与最高的OverspendAmount匹配,第二高的与第二高的匹配,以此类推,最终得到如下结果:
| Company Name | OverspendAmount | UnderspendAmount |
|---|---|---|
| Pepsi | 3600.00 | 2700.00 |
| Pepsi | 1400.00 | 1200.00 |
现有方案的性能问题
你当前使用collect的方案性能极差,核心原因是把全量数据拉取到Driver端单节点处理,完全没有利用Spark分布式计算的优势,数据量稍大就会触发Driver OOM、序列化反序列化开销过大等问题。
最优解决方案
直接使用Spark原生窗口函数+双表关联实现,全程分布式执行无数据落Driver的开销,是当前场景下性能最优的方案,代码示例如下:
from pyspark.sql import Window from pyspark.sql.functions import row_number, col # 步骤1:给df1每个公司下的UnderspendAmount按降序生成排名 w1 = Window.partitionBy("Company Name").orderBy(col("UnderspendAmount").desc()) df1_ranked = df1.withColumn("rank", row_number().over(w1)) # 步骤2:给df2每个公司下的OverspendAmount按降序生成排名 w2 = Window.partitionBy("Company Name").orderBy(col("OverspendAmount").desc()) df2_ranked = df2.withColumn("rank", row_number().over(w2)) # 步骤3:按公司+排名关联,取需要的字段即可 result = df2_ranked.join( df1_ranked, on=["Company Name", "rank"], how="inner" # 用inner仅保留两边都有的位次,用full可保留所有位次,缺失值补null ).select("Company Name", "OverspendAmount", "UnderspendAmount")
注:如果列名带空格运行报错,可以给列名加反引号包裹,比如
col("Company Name")
关于UDF+agg方案的说明
这个场景完全不需要用UDF+agg实现:
- 窗口函数是Spark Catalyst优化器原生支持的算子,执行效率比自定义UDF高3~10倍
- 若硬要用UDF+agg实现,需要先把同公司的所有金额collect成列表再按索引匹配,本质还是把聚合后的列表拉到内存计算,数据量大时同样有性能问题,稳定性远不如窗口函数方案
扩展性说明
当前窗口函数方案的扩展性极强,可支持后续功能扩展:
- 如果需要同金额同排名,把
row_number换成rank/dense_rank即可 - 如果要支持左/右/全连接保留某侧所有数据,修改join的how参数即可
- 如果后续要加其他匹配维度,只需要在窗口partitionBy和join的on条件里加对应字段即可
- 天然支持TB级大表分布式计算,不需要修改核心逻辑
内容的提问来源于stack exchange,提问作者amggg013
相关产品推荐
相关产品推荐

