如何高效移除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
相关产品推荐
相关产品推荐

