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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.29 08:44:58