PySpark大数据场景下Subtract操作过慢,求高效替代方案
高效实现优先级ID去重对比方案
数据示例
from pyspark.sql.types import StructType, StructField, StringType, IntegerType from pyspark.sql import SparkSession spark = SparkSession.builder.getOrCreate() schema = StructType([ StructField("ID", StringType(), True), StructField("Priority", IntegerType(), True) ]) data = spark.createDataFrame([ ("A", 1), ("B", 1), ("B", 2), ("C", 2), ("C", 3), ("D", 3) ], schema)
输入数据包含ID和Priority两列,核心记录如下:
- ID: A,仅出现于优先级1
- ID: B,出现于优先级1、2
- ID: C,出现于优先级2、3
- ID: D,仅出现于优先级3
需求说明
保留每个优先级中未在更低优先级出现过的ID:
- 优先级2仅保留未在优先级1出现的ID
- 优先级3仅保留未在优先级1、2出现的ID
- 以此类推
最终预期结果为:A(1)、B(1)、C(2)、D(3)
现有低效方案
当前方案通过逐轮迭代对比追加数据,随着new_data体积增大,多次subtract和union操作会引发频繁数据shuffle,大数据场景下性能极差:
步骤1:初始化基础数据
from pyspark.sql.functions import col new_data = data.filter(col('Priority') == 1)
步骤2:逐轮迭代处理
import pyspark.sql.functions as F for i in range(2, 4): x = data.filter(col('Priority') == i).select('ID') x = x.subtract(new_data.select('ID')) x = x.withColumn('Priority', F.lit(i)) new_data = new_data.union(x)
高效实现方案
利用Spark窗口函数一次性计算每个ID的首次出现的最低优先级,再筛选出符合条件的记录,全程仅需一次shuffle操作,性能远超迭代方案:
from pyspark.sql import Window import pyspark.sql.functions as F # 定义窗口:按ID分组 window_spec = Window.partitionBy("ID") # 计算每个ID的最小优先级(即首次出现的优先级) data_with_min_priority = data.withColumn( "min_priority", F.min("Priority").over(window_spec) ) # 筛选:仅保留当前优先级等于该ID最小优先级的记录 result = data_with_min_priority.filter( F.col("Priority") == F.col("min_priority") ).drop("min_priority") # 查看结果 result.show()
执行逻辑说明
- 窗口函数计算最小优先级:按ID分组后,用
min(Priority)得到每个ID首次出现的优先级 - 过滤符合条件的记录:如果当前记录的优先级等于该ID的最小优先级,说明此ID未在更低优先级出现过,保留该记录
- 结果匹配预期:最终输出与迭代方案的
new_data完全一致,但性能在大数据场景下提升显著
内容的提问来源于stack exchange,提问作者Vignesh A
相关产品推荐
相关产品推荐

