PySpark中dropDuplicates与Join操作结果异常问题咨询
问题分析与解决
你遇到的这个问题其实是PySpark中dropDuplicates()方法的默认行为导致的——当你只指定去重键(比如这里的ID)时,PySpark并不会保证保留重复组中的哪一行,它只会随机保留其中一行,这取决于底层的分区和执行计划,所以你看到的结果(保留了region_code=9的行)和预期(保留region_code=7)不一致是完全可能的。
为什么会这样?
dropDuplicates(["ID"])的逻辑是:对每个ID分组,然后从每组中挑选任意一行返回。它没有内置的规则来选择哪一行(比如最早出现的、某个字段最大/最小的),所以当同一个ID有多条记录时,最终保留的行是不确定的。你之前执行file_2.filter("ID == '1'").show()看到的是7,可能只是那次执行刚好保留了该行,但后续的去重操作(或者不同的执行上下文)可能选择了另一行。
解决方案:明确指定保留哪一行
如果你想确保保留特定的行(比如region_code最小的、或者最早加载的行),不能只依赖dropDuplicates(),而是需要用窗口函数来明确筛选规则。比如:
示例1:保留每个ID对应的最小region_code的行
from pyspark.sql.window import Window from pyspark.sql.functions import row_number # 定义窗口:按ID分组,按region_code升序排序 window_spec = Window.partitionBy("ID").orderBy("region_code") # 添加行号,每个ID组内行号为1的就是region_code最小的行 file_2 = file_2.withColumn("row_num", row_number().over(window_spec)) \ .filter("row_num == 1") \ .drop("row_num")
示例2:保留每个ID最早出现的行(如果有顺序标识字段)
如果你的数据有可以记录加载顺序的字段(比如load_timestamp),可以按该字段排序来保留最早的行:
window_spec = Window.partitionBy("ID").orderBy("load_timestamp") file_2 = file_2.withColumn("row_num", row_number().over(window_spec)) \ .filter("row_num == 1") \ .drop("row_num")
验证修改后的结果
当你用窗口函数明确筛选规则后,再执行join操作,就能确保每个ID对应的region_code是你预期的值了。比如用示例1的方法处理file_2后,ID=1的行就会保留region_code=7,关联后的结果也会完全符合你的预期。
内容的提问来源于stack exchange,提问作者default_settings
相关产品推荐
相关产品推荐

