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

在Databricks中用Python按目录查找基于主键的文件重复数据

基于主键(pk_id)分析嵌套目录CSV文件的重复与差异

实现步骤

  1. 加载嵌套目录文件并添加文件路径列
    要追踪每行数据的来源文件,需在加载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())
    
  2. 基于主键分析重复数据及对应文件
    仅用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.28 20:54:58