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

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()

核心调整逻辑

  1. 连接类型从inner改为full_outer,同时保留旧表独有的删除行、新表独有的新增行、两边共有的匹配行
  2. 新增status字段判断规则:
    • 旧表主键为空 → 行仅在新表存在,状态为added
    • 新表主键为空 → 行仅在旧表存在,状态为deleted
    • 匹配行存在更新列 → 状态为updated
    • 其余匹配行无变化 → 状态为unchanged
  3. 非主键列用coalesce处理,新增行取新表值、删除行取旧表值、匹配行取新表最新值
  4. 不等判断改用eqNullSafe,避免字段为空时判断逻辑出错

运行上述代码即可得到你期望的输出结果。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.27 10:54:04