You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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()

执行逻辑说明

  1. 窗口函数计算最小优先级:按ID分组后,用min(Priority)得到每个ID首次出现的优先级
  2. 过滤符合条件的记录:如果当前记录的优先级等于该ID的最小优先级,说明此ID未在更低优先级出现过,保留该记录
  3. 结果匹配预期:最终输出与迭代方案的new_data完全一致,但性能在大数据场景下提升显著

内容的提问来源于stack exchange,提问作者Vignesh A

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.07 12:35:17