PySpark drop_duplicates()执行后数据计数异常增加问题咨询
PySpark dropDuplicates() 去重后数据量反增的排查方案
问题背景
执行df.dropDuplicates()得到的df2计数为424527,原df计数仅为424510——去重后数据量反而上升。已排除以下低概率场景:
- 计数为近似值
- EMR Studio搭配12台m5.16xlarge集群的未知问题
- 源数据延迟加载导致的变更(连续执行计数语句验证过)
实用排查步骤
- 严格对齐计数逻辑
确保两次计数用的是完全一致的方法:必须都调用df.count(),禁止混用近似计数API。同时检查原df是否在计数前被隐式修改(比如触发过缓存、分区调整),导致两次计数的数据源本质不同。 - 核查dropDuplicates的参数与执行上下文
默认dropDuplicates()基于所有列去重,若指定了子集列(如dropDuplicates(["id", "create_time"])),要确认是否因业务逻辑误解导致预期偏差。另外,重点检查:原df是否在执行去重前被意外覆盖?比如是否写了df = df.some_transformation()但没保存到新变量,导致后续dropDuplicates()的输入已不是原始数据。 - 对比执行计划锁定数据源一致性
分别执行df.explain()和df2.explain(),对比两者的输入数据源是否完全一致。若原df依赖动态数据源(如Kafka流、未持久化临时视图),即使连续计数也可能因计划重跑导致数据变化——可以先执行df.cache().count()缓存原df,再基于缓存后的df做去重计数,彻底排除数据源动态变更的可能。 - 排查复杂数据类型与空值的干扰
PySpark对ArrayType、MapType等复杂类型的相等判断存在特殊逻辑,空值(null)的处理也可能导致重复行识别异常。可以先将复杂列转为字符串(如df.withColumn("array_str", concat_ws(",", col("array_col")))),再用转换后的列执行去重,验证是否是类型判断导致的问题。 - 缩小范围验证集群/分区问题
取原df的小样本(如df.limit(1000))执行去重,看是否仍出现数据量反增的情况。如果小样本正常,尝试强制单分区执行:df.repartition(1).dropDuplicates().count(),排查是否是多分区处理时的bug。还可以尝试重新分区原df(如df.repartition(200))后再去重,排除分区数据损坏的可能。 - 通过持久化隔离执行链路
对原df执行df.checkpoint()(需先设置checkpoint目录),将数据持久化到磁盘,再基于checkpoint后的df执行去重计数,隔离上游执行链路的潜在影响。
内容的提问来源于stack exchange,提问作者Gaurav Singhal
相关产品推荐
相关产品推荐

