Spark/PySpark中DataFrame连续两次show()结果不一致的问题
Spark中dropDuplicates去重后结果不稳定的问题与解决方案
问题场景
通过循环union多个Parquet文件生成DataFrame,数据中存在大量ID重复的行:
master_df = None for k, files in key_to_parquet_files.items(): df = read_parquet(*files) df = df.select(['ID']).withColumn('foo', F.lit(k)) if master_df is None: master_df = df else: master_df = master_df.union(df)
执行dropDuplicates(['ID'])去重后,两次查询同一ID集合的结果中,foo字段值不一致:
master_df = master_df.dropDuplicates(['ID']) print_ids_df = master_df.filter(F.col("ID").isin(*["A", "B", "C"])).orderBy(F.col('ID')) print_ids_df.show() print_ids_df = master_df.filter(F.col("ID").isin(*["A", "B", "C"])).orderBy(F.col('ID')) print_ids_df.show()
同一次运行中两次输出的foo值不同,所有针对ID的操作结果均不稳定。
问题根源
Spark的dropDuplicates是非确定性操作:当存在重复行时,Spark不会按照固定规则选择保留哪一行,而是随机选取重复组中的任意一行返回。这就导致即使是同一作业的多次执行(或同一作业内的不同阶段),都可能得到不同的结果。
确定性去重方案
要得到稳定可预测的去重结果,需要明确指定重复组中保留行的规则。使用窗口函数可以实现这一点:
from pyspark.sql import Window import pyspark.sql.functions as F # 按ID分组,按foo字段升序排序(空值后置) master_df_window = Window.partitionBy(['ID']).orderBy(F.col('foo').asc_nulls_last()) # 添加行号,只保留每组的第一行 master_df = master_df.withColumn('rank', F.row_number().over(master_df_window)) \ .filter(F.col('rank') == 1) \ .drop('rank')
这段代码需要在dropDuplicates之前执行,通过partitionBy('ID')将相同ID的行归为一组,orderBy定义了组内行的排序规则(这里选择保留foo值最小的行),最后通过row_number()和过滤条件固定保留每组的第一行,确保结果完全可预测。
内容的提问来源于stack exchange,提问作者rouble
相关产品推荐
相关产品推荐

