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

PySpark连续调用dropDuplicates后记录数异常增加问题求助

问题分析与解决方案

首先明确:正常情况下,dropDuplicates的输出记录数必然≤输入记录数,你看到第二次去重后记录数增加的情况,肯定是某个环节出现了异常,以下是可能的原因及排查方案:

1. 日志输出逻辑错误

优先检查日志打印的准确性,确认两次logger.info对应的确实是两次去重后的df.count()结果。比如是否不小心将第一次的统计值赋值给了第二次的变量,或者日志语句的顺序写反了。

2. Spark惰性求值导致的执行链异常

Spark是惰性计算模型,只有遇到count()这类行动操作时才会触发计算。如果两次去重后的df没有被持久化,第二次count()可能会重新计算整个执行链(原始数据→第一次去重→第二次去重),虽然理论上结果仍应≤第一次去重后的记录数,但可以通过缓存验证:

df = spark.read.format("parquet").load(my_table_path)    
df = df.dropDuplicates(("Column_1", "Column_2"))
df.cache()  # 缓存第一次去重后的结果
logger.info("No. of records after first drop: {}".format(df.count()))
df = df.dropDuplicates(("Column_1", "Column_3"))
logger.info("No. of records after second drop: {}".format(df.count()))
df.unpersist()

缓存后第二次去重会基于已缓存的数据集执行,避免重复读取原始数据。

3. 数据中null值的影响

Spark中null值与任何值(包括另一个null)都不相等。如果Column_3存在大量null,第一次去重后的数据集里,同一个Column_1对应的Column_3=null的记录会被全部保留(因为(Column_1, null)被视为不重复组合),但这只会让第二次去重后的记录数等于第一次的,不会增加。可以检查Column_3的null占比:

from pyspark.sql.functions import col
null_count = df.filter(col("Column_3").isNull()).count()
logger.info("Null count in Column_3: {}".format(null_count))

4. 用分组聚合替代dropDuplicates验证逻辑

可以用更直观的分组聚合方式替代dropDuplicates,验证去重结果是否符合预期:

df = spark.read.format("parquet").load(my_table_path)    
# 第一次按(Column_1, Column_2)去重,保留任意一条记录的其他列
df = df.groupBy("Column_1", "Column_2").agg(
    col("Column_3").alias("Column_3"),
    col("Column_4").alias("Column_4")
)
logger.info("No. of records after first drop: {}".format(df.count()))
# 第二次按(Column_1, Column_3)去重
df = df.groupBy("Column_1", "Column_3").agg(
    col("Column_2").alias("Column_2"),
    col("Column_4").alias("Column_4")
)
logger.info("No. of records after second drop: {}".format(df.count()))

5. 数据倾斜导致的计算异常

如果数据集存在严重的数据倾斜(比如某个Column_1值对应数百万条记录),可能导致count计算出现异常。可以尝试重新分区后再执行去重:

df = spark.read.format("parquet").load(my_table_path)
df = df.repartition(200)  # 根据集群资源调整分区数
df = df.dropDuplicates(("Column_1", "Column_2"))
logger.info("No. of records after first drop: {}".format(df.count()))
df = df.dropDuplicates(("Column_1", "Column_3"))
logger.info("No. of records after second drop: {}".format(df.count()))

内容的提问来源于stack exchange,提问作者Aart

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.19 20:40:50