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

如何高效移除PySpark DataFrame数百万行数据中的重复项

高效解决PySpark复杂Schema展平后的数据去重问题

我有一段可正常运行的PySpark代码,用来展平包含ArrayType和StructType的复杂Schema,之后筛选指定条件的数据,但生成的df2存在大量重复项。试过dropDuplicates()、distinct()、重分区、缓存这些去重方法,都耗时20分钟以上还没成功,求高效准确的去重方案。

核心优化思路

问题根源是过早全量展平Schema,导致数据爆炸式膨胀,后续去重的计算成本被无限放大。正确的思路是先筛选再展平,只处理业务需要的字段,从源头减少数据规模。

具体优化步骤

  • 提前过滤核心条件:先根据cvrNummer和vaerdi="REVISION"筛选数据,只保留目标行,砍掉绝大多数无关数据。
  • 按需展平嵌套字段:不需要把整个Schema拆平,只针对查询用到的嵌套数组、结构体进行拆解,避免不必要的行膨胀。
  • 在最小数据集上去重:等数据量压缩到最小后再执行去重操作,而非等数据膨胀后再处理。
  • 利用Spark惰性求值:让过滤操作优先执行,减少后续所有步骤的数据处理量。

修改后的代码

from pyspark.sql.types import *
from pyspark.sql.functions import *

# 仅针对业务需要的字段进行部分展平,避免全表无意义膨胀
def flatten_required_fields(df):
    # 第一步:先过滤核心条件,大幅缩小数据范围
    filtered_df = df.filter(
        col("_source.Vrvirksomhed.cvrNummer") == 10088437
    ).filter(
        exists(
            "_source.Vrvirksomhed.deltagerRelation.organisationer.medlemsData.attributter.vaerdier",
            lambda x: x.vaerdi == "REVISION"
        )
    )
    
    # 第二步:仅拆解需要的数组字段
    exploded_df = filtered_df.withColumn(
        "vaerdier",
        explode_outer("_source.Vrvirksomhed.deltagerRelation.organisationer.medlemsData.attributter.vaerdier")
    )
    
    # 第三步:筛选数组中符合条件的元素,并提取目标字段
    final_df = exploded_df.filter(col("vaerdier.vaerdi") == "REVISION").select(
        col("_source.Vrvirksomhed.cvrNummer").alias("_source_Vrvirksomhed_cvrNummer"),
        col("vaerdier.vaerdi").alias("_source_Vrvirksomhed_deltagerRelation_organisationer_medlemsData_attributter_vaerdier_vaerdi"),
        col("vaerdier.periode.gyldigFra").alias("_source_Vrvirksomhed_deltagerRelation_organisationer_medlemsData_attributter_vaerdier_periode_gyldigFra"),
        col("vaerdier.periode.gyldigTil").alias("_source_Vrvirksomhed_deltagerRelation_organisationer_medlemsData_attributter_vaerdier_periode_gyldigTil")
    )
    
    # 第四步:在最小数据集上执行去重
    return final_df.dropDuplicates()

df2 = flatten_required_fields(df)

额外性能优化建议

  • 分区利用:如果原始数据按cvrNummer分区存储,直接通过分区过滤可以跳过大量无关数据的扫描。
  • 缓存策略:在筛选后的filtered_df上执行cache()或persist(StorageLevel.MEMORY_AND_DISK),避免重复计算嵌套字段的拆解逻辑。
  • 执行计划校验:运行df2.explain()查看执行计划,确保过滤操作被下推到数据源层(如Parquet/ORC支持谓词下推),进一步提升效率。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.13 13:51:15