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

基于Primary Key对比同Schema DataFrame差异并输出结果的技术问询

我来帮你搞定这个DataFrame对比的需求,下面是具体的实现思路和代码示例(以PySpark为例),完全贴合你提到的要求:

核心实现思路
  • 用full outer join关联两个DataFrame:这是实现“保留所有主键记录(包括仅在单个表存在的)”的关键,能确保不会漏掉任何一条主键数据
  • 动态遍历非主键列:对每一列做差异判断,还要特殊处理NULL值的情况(Spark中NULL和任何值对比结果都是NULL,必须单独判断)
  • 分两类输出结果:一类是仅在单个DataFrame中存在的主键记录,另一类是主键存在但列值有差异的记录,每类结果都做清晰格式化
代码实现步骤

假设我们的主键列名为primary_key,两个待对比的DataFrame分别是df1和df2,且二者Schema完全一致。

第一步:重命名列并执行Full Outer Join

先给两个DataFrame的列名加上后缀区分来源,再执行关联:

from pyspark.sql import functions as F
from pyspark.sql.types import *
from functools import reduce

# 给df1和df2的列名分别加后缀,避免关联后列名冲突
df1_renamed = df1.select([F.col(c).alias(f"{c}_df1") for c in df1.columns])
df2_renamed = df2.select([F.col(c).alias(f"{c}_df2") for c in df2.columns])

# 执行full outer join,关联主键列
joined_df = df1_renamed.join(
    df2_renamed,
    df1_renamed["primary_key_df1"] == df2_renamed["primary_key_df2"],
    how="full_outer"
)

第二步:提取仅在单个DataFrame存在的主键记录

通过判断主键列是否为NULL,筛选出仅在某一个表中存在的记录:

# 筛选仅在df1中存在的记录
only_in_df1 = joined_df.filter(F.col("primary_key_df2").isNull()).select(
    F.col("primary_key_df1").alias("primary_key"),
    F.lit("仅存在于DataFrame1").alias("status")
)

# 筛选仅在df2中存在的记录
only_in_df2 = joined_df.filter(F.col("primary_key_df1").isNull()).select(
    F.col("primary_key_df2").alias("primary_key"),
    F.lit("仅存在于DataFrame2").alias("status")
)

# 合并两类单独存在的记录
single_source_records = only_in_df1.union(only_in_df2)

第三步:提取主键存在但列值有差异的记录

动态遍历所有非主键列,生成差异判断逻辑,最后整理成易读的差异详情:

# 获取所有非主键列名
non_key_cols = [col for col in df1.columns if col != "primary_key"]

# 为每个非主键列生成差异判断表达式(包含NULL值处理)
diff_conditions = []
for col_name in non_key_cols:
    col_df1 = f"{col_name}_df1"
    col_df2 = f"{col_name}_df2"
    # 差异判断:要么值不相等,要么一个是NULL另一个不是
    condition = (F.col(col_df1) != F.col(col_df2)) | (F.col(col_df1).isNull() != F.col(col_df2).isNull())
    diff_conditions.append(condition.alias(f"{col_name}_has_diff"))

# 构建包含主键、差异标记、原始值的基础DataFrame
diff_base_df = joined_df.select(
    F.coalesce(F.col("primary_key_df1"), F.col("primary_key_df2")).alias("primary_key"),
    *diff_conditions,
    *[F.col(f"{c}_df1").alias(f"{c}_df1") for c in non_key_cols],
    *[F.col(f"{c}_df2").alias(f"{c}_df2") for c in non_key_cols]
)

# 筛选出至少有一个列存在差异的记录
diff_records = diff_base_df.filter(
    reduce(lambda a, b: a | b, [F.col(f"{c}_has_diff") for c in non_key_cols])
)

# 定义UDF来整理差异详情,把每个差异列的信息打包成结构化数据
def extract_diff_details(row):
    diff_list = []
    for col_name in non_key_cols:
        if row[f"{col_name}_has_diff"]:
            val1 = row[f"{col_name}_df1"]
            val2 = row[f"{col_name}_df2"]
            diff_list.append({
                "column_name": col_name,
                "df1_value": str(val1) if val1 is not None else "NULL",
                "df2_value": str(val2) if val2 is not None else "NULL"
            })
    return diff_list

# 定义UDF的返回Schema
diff_schema = ArrayType(StructType([
    StructField("column_name", StringType()),
    StructField("df1_value", StringType()),
    StructField("df2_value", StringType())
]))

diff_udf = F.udf(extract_diff_details, diff_schema)

# 生成最终的差异结果
final_diff_records = diff_records.select(
    F.col("primary_key"),
    diff_udf(F.struct([F.col(c) for c in diff_records.columns])).alias("differences")
)

第四步:输出结果

最后可以分别打印两类结果,查看对比情况:

print("=== 仅存在于单个DataFrame的主键记录 ===")
single_source_records.show(truncate=False)

print("\n=== 存在列值差异的主键记录及详情 ===")
final_diff_records.show(truncate=False)
补充说明
  • 如果是使用Scala版本的Spark,思路完全一致,只是语法上需要做对应调整(比如用reduceLeft代替Python的reduce,用udf的定义方式不同等)
  • 处理NULL值是关键:如果忽略NULL的情况,会漏掉很多真实存在的差异(比如一个表中某列是NULL,另一个表中是有效值)

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.25 03:52:51