如何将Pandas中校验数据差异的merge逻辑转换为PySpark实现?
PySpark实现DataFrame基准值校验逻辑
需求说明
现有两个结构相同的DataFrame,各含3列。需要校验每组col_1、col_2对应的col_to_check值是否与基准DataFrame一致,输出存在差异的行。已通过Pandas的merge+query实现,现需用PySpark完成相同逻辑。
基准DataFrame(input_benchmark)
| col_1 | col_2 | col_to_check |
|---|---|---|
| girl | 12 | Primary |
| boy | 14 | Secondary |
| baby | 1 | Nursery |
| girl_1 | 10 | Secondary |
| girl_2 | 10 | Secondary |
待校验DataFrame(input_df)
| col_1 | col_2 | col_to_check |
|---|---|---|
| girl | 12 | Primary |
| boy | 14 | Secondary |
| baby | 1 | Secondary |
| toddler | 3 | Kindergarten |
| girl_1 | 10 | null |
| girl_2 | 10 | null |
已实现的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
相关产品推荐
相关产品推荐

