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

在Databricks中执行Upsert后,对比两个DataFrame的变更内容

在Databricks中对比Upsert前后DataFrame的变更情况

假设你的主键是id,可以通过全外连接+字段对比的方式快速得到Upsert前后的数据变更情况,以下是具体实现步骤:

1. 模拟示例数据(对应你的df1和df2)

# 原始DataFrame df1
df1 = spark.createDataFrame(
    [(1, "Alice", 25), (2, "Bob", 30), (3, "Charlie", 35)],
    ["id", "name", "age"]
)

# Upsert后的DataFrame df2
df2 = spark.createDataFrame(
    [(1, "Alice", 26), (2, "Bob", 30), (4, "David", 28)],
    ["id", "name", "age"]
)

2. 全外连接两个DataFrame

通过全外连接覆盖所有id的存在情况(仅在df1、仅在df2、同时存在于两个DF):

joined_df = df1.alias("old").join(df2.alias("new"), on="id", how="full_outer")

3. 标记变更类型

根据id的存在情况和字段值差异,标记每条数据的变更类型:

from pyspark.sql.functions import when, concat, lit, concat_ws

change_type_df = joined_df.withColumn(
    "change_type",
    when(col("old.id").isNull(), "新增")
    .when(col("new.id").isNull(), "删除")
    .when(
        (col("old.name") != col("new.name")) | (col("old.age") != col("new.age")),
        "修改"
    )
    .otherwise("无变更")
)

4. 生成变更详情

列出具体修改的字段及前后值(用Spark内置函数实现,避免UDF提升性能):

result_df = change_type_df.withColumn(
    "change_details",
    concat_ws(
        "; ",
        when(col("old.name") != col("new.name"), concat(lit("name: "), col("old.name"), lit(" → "), col("new.name"))),
        when(col("old.age") != col("new.age"), concat(lit("age: "), col("old.age"), lit(" → "), col("new.age")))
    )
)

5. 整理并展示结果

选择需要的列,按主键排序后输出:

final_result = result_df.select(
    "id",
    col("old.name").alias("old_name"),
    col("new.name").alias("new_name"),
    col("old.age").alias("old_age"),
    col("new.age").alias("new_age"),
    "change_type",
    "change_details"
).orderBy("id")

final_result.show(truncate=False)

输出结果示例:

+---+---------+---------+-------+-------+-----------+------------------+
|id |old_name |new_name |old_age|new_age|change_type|change_details    |
+---+---------+---------+-------+-------+-----------+------------------+
|1  |Alice    |Alice    |25     |26     |修改        |age: 25 → 26      |
|2  |Bob      |Bob      |30     |30     |无变更      |null              |
|3  |Charlie  |null     |35     |null   |删除        |null              |
|4  |null     |David    |null   |28     |新增        |null              |
+---+---------+---------+-------+-------+-----------+------------------+

如果你的DataFrame有更多字段,只需要在change_type的判断条件和change_details的拼接逻辑中添加对应字段即可。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.04 21:25:13