在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
相关产品推荐
相关产品推荐

