在Databricks中用Python按目录查找基于主键的文件重复数据
基于主键(pk_id)分析嵌套目录CSV文件的重复与差异
实现步骤
加载嵌套目录文件并添加文件路径列
要追踪每行数据的来源文件,需在加载DataFrame时通过input_file_name()函数新增文件路径字段,这样后续分析重复主键时能直接对应到具体文件。示例代码:
from pyspark.sql.functions import input_file_name # 加载嵌套目录下所有CSV文件,同时添加文件路径列 df = spark.read.format('csv').options(header=True).load('/mnt/home/Archive/abc/*/*/*/*.csv') \ .withColumn("file_path", input_file_name())基于主键分析重复数据及对应文件
仅用dropDuplicates()只能得到去重后的数据集,无法直接定位哪些文件存在重复主键。需先筛选出重复的主键,再关联对应的文件路径和行数据:# 找出所有出现多次的pk_id duplicate_pks = df.groupBy("pk_id").count().filter("count > 1").select("pk_id") # 关联原数据集,获取重复主键对应的所有行及文件路径 duplicate_rows_with_files = df.join(duplicate_pks, on="pk_id", how="inner") # 按pk_id分组,查看每个重复主键对应的文件集合和所有行数据(便于对比差异) duplicate_details = duplicate_rows_with_files.groupBy("pk_id").agg( collect_set("file_path").alias("related_files"), collect_list(struct(*df.columns)).alias("all_rows") ) # 展示分析结果 duplicate_details.show(truncate=False)
关键说明
input_file_name():Spark内置函数,返回当前行数据所属的文件路径,用于数据溯源。- 通过分组统计主键出现次数筛选重复项,再关联原数据可完整呈现每个重复主键对应的文件和行内容,方便对比行数据差异。
内容的提问来源于stack exchange,提问作者user2181700
相关产品推荐
相关产品推荐

