以下代码在 Databricks (Serverless) 中运行,通过贪心后向消除算法找出表的最短唯一键。若整表存在完全重复的行,则生成一条 SELECT 语句用于查看其中一组重复记录。
frompyspark.sqlimportfunctionsasFfrompyspark.sqlimportSparkSession# 假设 spark 已由 Databricks 环境提供# 输入表名,可通过 widget 或直接赋值table_name="your_catalog.your_schema.your_table"# 请替换为实际表名# 1. 读取表并获取总行数df=spark.table(table_name)total_rows=df.count()all_cols=df.columnsprint(f"表总行数:{total_rows}")print(f"全部字段:{all_cols}")# 2. 检查整表是否完全去重(所有字段组合)full_distinct=df.select(all_cols).distinct().count()iffull_distinct<total_rows:print("表在所有字段上存在完全重复的行,正在生成筛选一条重复记录的 SELECT 语句...")# 取一组重复记录:按所有字段分组,取 count>1 的第一组dup_group=(df.groupBy(all_cols).agg(F.count("*").alias("cnt")).filter(F.col("cnt")>1).select(all_cols).limit(1).collect())ifdup_group:row=dup_group[0]# 构建 WHERE 条件,注意处理 NULL 和数据类型conditions=[]forcol_nameinall_cols:val=row[col_name]ifvalisNone:conditions.append(f"{col_name}IS NULL")else:# 字符串类型需要加引号,简单起见使用 repr() 转为字面量conditions.append(f"{col_name}={repr(val)}")where_clause=" AND ".join(conditions)select_sql=f"SELECT * FROM{table_name}WHERE{where_clause};"print("筛选一组重复记录的 SQL 语句:")print(select_sql)else:print("未找到重复记录,请检查数据。")else:# 3. 贪心后向消除寻找最短唯一键current_cols=all_cols.copy()# 当前候选列# 不断尝试移除一列,直到无法再移除changed=Truewhilechanged:changed=False# 遍历当前候选列,寻找可以移除的列forcolincurrent_cols:test_cols=[cforcincurrent_colsifc!=col]# 计算 test_cols 的去重行数test_distinct=df.select(test_cols).distinct().count()iftest_distinct==total_rows:# 移除该列后仍能唯一标识,执行移除并跳出本次循环current_cols=test_cols changed=Truebreak# 修改了列组合,重新开始检查print(f"找到的最短唯一键:{current_cols}")print(f"该组合字段数:{len(current_cols)}")# 可选:展示该唯一键的去重行数(应等于总行数)final_distinct=df.select(current_cols).distinct().count()print(f"唯一键去重行数:{final_distinct}(总行数:{total_rows})")说明:
- 不持久化:代码没有使用 persist()、cache() 或 checkpoint(),完全依赖 Spark 的惰性计算,每次 count() 都会触发重算表,符合 Serverless 限制。
- 算法逻辑:从全部字段开始,贪心地尝试删减字段,只要删减后去重行数仍等于表总行数,就永久删除该字段,直到无法删除为止。
- 重复处理:若所有字段去重后仍小于总行数,则通过分组找出一组重复值,并用这些值拼接成 SELECT * FROM 表 WHERE 条件 的语句,可直接运行查看重复行。
- NULL 安全:生成条件时对 None 使用 IS NULL,其他值使用 repr() 保证字符串正确转义。
可将 table_name 替换为实际的三级表名(如 hive_metastore.default.my_table)后直接在 Databricks Notebook 中运行。