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
相关产品推荐
相关产品推荐

