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

如何将Pandas中校验数据差异的merge逻辑转换为PySpark实现?

PySpark实现DataFrame基准值校验逻辑

需求说明

现有两个结构相同的DataFrame,各含3列。需要校验每组col_1、col_2对应的col_to_check值是否与基准DataFrame一致,输出存在差异的行。已通过Pandas的merge+query实现,现需用PySpark完成相同逻辑。

基准DataFrame(input_benchmark)

col_1col_2col_to_check
girl12Primary
boy14Secondary
baby1Nursery
girl_110Secondary
girl_210Secondary

待校验DataFrame(input_df)

col_1col_2col_to_check
girl12Primary
boy14Secondary
baby1Secondary
toddler3Kindergarten
girl_110null
girl_210null

已实现的Pandas代码

def check_func(input_benchmark, input_df):
    df_new = input_df.merge(input_benchmark, on=['col_1', 'col_2'], suffixes=(None, '_actual')).query('col_to_check != col_to_check_actual')
    return df_new

PySpark实现方案

PySpark中通过join替代Pandas的merge,需注意处理null值的比较逻辑(PySpark中null与任何值比较返回null,需单独判断):

校验函数实现

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

def check_func_spark(input_benchmark, input_df):
    # 内连接两个DataFrame,关联col_1和col_2,并重命名基准列
    joined_df = input_df.join(
        input_benchmark,
        on=["col_1", "col_2"],
        how="inner"
    ).select(
        input_df["col_1"],
        input_df["col_2"],
        input_df["col_to_check"],
        input_benchmark["col_to_check"].alias("col_to_check_actual")
    )
    # 过滤值不一致的行,包含一方为null的场景
    diff_df = joined_df.filter(
        (col("col_to_check") != col("col_to_check_actual")) |
        (col("col_to_check").isNull() & col("col_to_check_actual").isNotNull()) |
        (col("col_to_check").isNotNull() & col("col_to_check_actual").isNull())
    )
    return diff_df

测试代码

# 初始化SparkSession
spark = SparkSession.builder.appName("DataFrameCheck").getOrCreate()

# 创建基准DataFrame
benchmark_data = [
    ("girl", 12, "Primary"),
    ("boy", 14, "Secondary"),
    ("baby", 1, "Nursery"),
    ("girl_1", 10, "Secondary"),
    ("girl_2", 10, "Secondary")
]
input_benchmark = spark.createDataFrame(benchmark_data, ["col_1", "col_2", "col_to_check"])

# 创建待校验DataFrame
check_data = [
    ("girl", 12, "Primary"),
    ("boy", 14, "Secondary"),
    ("baby", 1, "Secondary"),
    ("toddler", 3, "Kindergarten"),
    ("girl_1", 10, None),
    ("girl_2", 10, None)
]
input_df = spark.createDataFrame(check_data, ["col_1", "col_2", "col_to_check"])

# 执行校验并展示结果
result_df = check_func_spark(input_benchmark, input_df)
result_df.show()

输出结果

+-------+-----+-------------+-------------------+
|  col_1|col_2|col_to_check|col_to_check_actual|
+-------+-----+-------------+-------------------+
|    baby|    1|    Secondary|            Nursery|
| girl_1|   10|         null|          Secondary|
| girl_2|   10|         null|          Secondary|
+-------+-----+-------------+-------------------+

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.16 20:01:13