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

如何逐列比较两个PySpark DataFrames并将结果并排追加?

PySpark 逐列对比两个DataFrame并生成结果

我有两个PySpark DataFrame,需要逐列比较它们的值,并将每列的原始值和对比结果并列展示,最终得到目标DataFrame。

原始数据

DF1

Claim_numberClaim_Status
1001Closed
1002In Progress
1003open

DF2

Claim_numberClaim_Status
1001Closed
1002open
1004In Progress

期望输出DF3

Claim_number_DF1Claim_number_DF2Comparison_of_Claim_numberClaim_status_DF1Claim_status_DF2Comparison_of_Claim_Status
10011001TRUEClosedClosedTRUE
10021002TRUEIn ProgressopenFALSE
10031004FALSEopenIn ProgressFALSE

实现代码

from pyspark.sql import SparkSession
from pyspark.sql.functions import col, monotonically_increasing_id

# 初始化Spark会话
spark = SparkSession.builder.appName("DFColumnComparison").getOrCreate()

# 构建原始DataFrame
data_df1 = [(1001, "Closed"), (1002, "In Progress"), (1003, "open")]
df1 = spark.createDataFrame(data_df1, ["Claim_number", "Claim_Status"])

data_df2 = [(1001, "Closed"), (1002, "open"), (1004, "In Progress")]
df2 = spark.createDataFrame(data_df2, ["Claim_number", "Claim_Status"])

# 添加行号,确保按原始行顺序关联
df1_with_rowid = df1.withColumn("row_id", monotonically_increasing_id())
df2_with_rowid = df2.withColumn("row_id", monotonically_increasing_id())

# 重命名列,区分来源
df1_renamed = df1_with_rowid.select(
    "row_id",
    col("Claim_number").alias("Claim_number_DF1"),
    col("Claim_Status").alias("Claim_status_DF1")
)

df2_renamed = df2_with_rowid.select(
    "row_id",
    col("Claim_number").alias("Claim_number_DF2"),
    col("Claim_Status").alias("Claim_status_DF2")
)

# 关联两个DataFrame
joined_df = df1_renamed.join(df2_renamed, on="row_id", how="inner")

# 生成对比结果列并调整列顺序
final_df = joined_df.select(
    "Claim_number_DF1",
    "Claim_number_DF2",
    (col("Claim_number_DF1") == col("Claim_number_DF2")).alias("Comparison_of_Claim_number"),
    "Claim_status_DF1",
    "Claim_status_DF2",
    (col("Claim_status_DF1") == col("Claim_status_DF2")).alias("Comparison_of_Claim_Status")
).orderBy("row_id").drop("row_id")

# 查看结果
final_df.show()

代码说明

  1. 行号关联:使用monotonically_increasing_id()生成唯一行标识,保证两个DataFrame按原始行顺序一一对应,这是实现逐行逐列对比的关键。
  2. 列重命名:给两个DataFrame的列添加_DF1/_DF2后缀,避免关联时列名冲突,同时清晰区分数据来源。
  3. 对比逻辑:直接用==运算符比较对应列的值,生成布尔类型的对比结果列,符合期望输出的TRUE/FALSE格式。
  4. 列顺序调整:按照目标DataFrame的列顺序重新排列,并移除临时的行号列。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.16 05:03:10