如何逐列比较两个PySpark DataFrames并将结果并排追加?
PySpark 逐列对比两个DataFrame并生成结果
我有两个PySpark DataFrame,需要逐列比较它们的值,并将每列的原始值和对比结果并列展示,最终得到目标DataFrame。
原始数据
DF1
| Claim_number | Claim_Status |
|---|---|
| 1001 | Closed |
| 1002 | In Progress |
| 1003 | open |
DF2
| Claim_number | Claim_Status |
|---|---|
| 1001 | Closed |
| 1002 | open |
| 1004 | In Progress |
期望输出DF3
| Claim_number_DF1 | Claim_number_DF2 | Comparison_of_Claim_number | Claim_status_DF1 | Claim_status_DF2 | Comparison_of_Claim_Status |
|---|---|---|---|---|---|
| 1001 | 1001 | TRUE | Closed | Closed | TRUE |
| 1002 | 1002 | TRUE | In Progress | open | FALSE |
| 1003 | 1004 | FALSE | open | In Progress | FALSE |
实现代码
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()
代码说明
- 行号关联:使用
monotonically_increasing_id()生成唯一行标识,保证两个DataFrame按原始行顺序一一对应,这是实现逐行逐列对比的关键。 - 列重命名:给两个DataFrame的列添加
_DF1/_DF2后缀,避免关联时列名冲突,同时清晰区分数据来源。 - 对比逻辑:直接用
==运算符比较对应列的值,生成布尔类型的对比结果列,符合期望输出的TRUE/FALSE格式。 - 列顺序调整:按照目标DataFrame的列顺序重新排列,并移除临时的行号列。
内容的提问来源于stack exchange,提问作者Srinivasan
相关产品推荐
相关产品推荐

