Pyspark对比两个Dataframe实现增删改行及更新字段识别
PySpark对比两张表识别增删改行的代码调整
以下是修改后的完整可运行代码:
from pyspark.sql import SparkSession from pyspark.sql.functions import col, array, when, array_remove, lit, coalesce, size from pyspark.sql.types import StructType, StructField, StringType, IntegerType # 环境已有SparkSession可省略初始化步骤 spark = SparkSession.builder.appName("table_compare").getOrCreate() data1 = [("James","rob","Smith","36636","M",3000), ("Michael","Rose","jim","40288","M",4000), ("Robert","dunkin","Williams","42114","M",4000), ("Maria","Anne","Jones","39192","F",4000), ("Jen","Mary","Brown","60563","F",-1) ] data2 = [("James","rob","Smith","36636","M",3000), ("Robert","dunkin","Williams","42114","M",2000), ("Maria","Anne","Jones","72712","F",3000), ("Yesh","Reddy","Brown","75234","M",3000), ("Jen","Mary","Brown","60563","F",-1) ] schema = StructType([ StructField("firstname",StringType(),True), StructField("middlename",StringType(),True), StructField("lastname",StringType(),True), StructField("id", StringType(), True), StructField("gender", StringType(), True), StructField("salary", IntegerType(), True) ]) df1 = spark.createDataFrame(data=data1,schema=schema) df2 = spark.createDataFrame(data=data2,schema=schema) # 定义主键、非主键列方便复用 primary_keys = ["firstname","middlename","lastname"] non_primary_cols = [c for c in df1.columns if c not in primary_keys] # 改用eqNullSafe判断不等,避免空值场景下判断异常 conditions_ = [when(~df1[c].eqNullSafe(df2[c]), lit(c)).otherwise("") for c in non_primary_cols] # 定义状态判断逻辑 status_col = when(df1[primary_keys[0]].isNull(), lit("added"))\ .when(df2[primary_keys[0]].isNull(), lit("deleted"))\ .when(size(array_remove(array(*conditions_), "")) > 0, lit("updated"))\ .otherwise(lit("unchanged")) select_expr = [ # 主键取两边非空值即可 *[coalesce(df1[k], df2[k]).alias(k) for k in primary_keys], # 非主键列优先取新表值,删除行 fallback 到旧表值 *[coalesce(df2[c], df1[c]).alias(c) for c in non_primary_cols], array_remove(array(*conditions_), "").alias("updated_columns"), status_col.alias("status") ] # 改用全外连接保留所有行 df1.join(df2, on=primary_keys, how="full_outer").select(*select_expr).show()
核心调整逻辑
- 连接类型从
inner改为full_outer,同时保留旧表独有的删除行、新表独有的新增行、两边共有的匹配行 - 新增
status字段判断规则:- 旧表主键为空 → 行仅在新表存在,状态为
added - 新表主键为空 → 行仅在旧表存在,状态为
deleted - 匹配行存在更新列 → 状态为
updated - 其余匹配行无变化 → 状态为
unchanged
- 旧表主键为空 → 行仅在新表存在,状态为
- 非主键列用
coalesce处理,新增行取新表值、删除行取旧表值、匹配行取新表最新值 - 不等判断改用
eqNullSafe,避免字段为空时判断逻辑出错
运行上述代码即可得到你期望的输出结果。
内容的提问来源于stack exchange,提问作者yAsH
相关产品推荐
相关产品推荐

