PySpark中对比同结构DataFrame,提取不匹配列及值
PySpark 对比两个同Schema DataFrame的不匹配数据
假设你有两个Schema完全一致的DataFrame df1 和 df2,需要以key_col为键执行内连接,提取每个键对应的不匹配列名及两个DataFrame的对应值,以下是实现代码:
实现步骤
- 提取除连接键外的所有列名
- 内连接两个DataFrame,为第二个DF的列添加后缀区分
- 生成列不匹配的判断条件,过滤出存在不匹配的行
- 使用
stack函数将列级不匹配转换为行级记录,整理成目标输出格式
完整代码
from pyspark.sql import SparkSession from pyspark.sql.functions import col, expr, stack # 初始化SparkSession(未初始化时执行) spark = SparkSession.builder.appName("DataFrameMismatchCheck").getOrCreate() # 示例数据(替换为你的实际DataFrame即可) data1 = [("key1", "val1", 100), ("key2", "val2", 200), ("key3", "val3", 300)] data2 = [("key1", "val1", 100), ("key2", "val2_updated", 200), ("key3", "val3", 301)] schema = ["key_col", "str_col", "num_col"] df1 = spark.createDataFrame(data1, schema=schema) df2 = spark.createDataFrame(data2, schema=schema) # 1. 获取除key_col外的所有列 non_key_cols = [col_name for col_name in df1.columns if col_name != "key_col"] # 2. 内连接并给两个DF的列加前缀区分 joined_df = df1.join(df2, on="key_col", how="inner")\ .select( "key_col", *[col(f"{col_name}").alias(f"{col_name}_df1") for col_name in non_key_cols], *[col(f"{col_name}").alias(f"{col_name}_df2") for col_name in non_key_cols] ) # 3. 生成不匹配条件,过滤出至少有一个列不匹配的行 # 若要处理NULL值,可将条件改为 f"{col_name}_df1 <=> {col_name}_df2 == False" mismatch_conditions = [expr(f"{col_name}_df1 != {col_name}_df2") for col_name in non_key_cols] filtered_df = joined_df.filter(expr(" OR ".join([cond.desc for cond in mismatch_conditions]))) # 4. 将列转换为行,提取不匹配的列名和对应值 stack_expr = f"stack({len(non_key_cols)}, {', '.join([f'{repr(col_name)}, {col_name}_df1, {col_name}_df2' for col_name in non_key_cols])}) as (mismatched_col, df1_value, df2_value)" result_df = filtered_df.select("key_col", expr(stack_expr))\ .filter(col("df1_value") != col("df2_value"))\ .orderBy("key_col", "mismatched_col") # 查看结果 result_df.show(truncate=False)
输出示例
+-------+-------------+----------+----------+ |key_col|mismatched_col|df1_value |df2_value | +-------+-------------+----------+----------+ |key2 |str_col |val2 |val2_updated| |key3 |num_col |300 |301 | +-------+-------------+----------+----------+
补充说明
- 如果列中存在
NULL值,建议使用Spark的安全相等运算符<=>来判断不匹配,避免NULL值被忽略,对应条件修改为expr(f"{col_name}_df1 <=> {col_name}_df2") == False - 若需要保留所有
key_col的记录(包括完全匹配的),可以跳过filtered_df的过滤步骤,直接执行stack后过滤不匹配行即可
内容的提问来源于stack exchange,提问作者Rahul Kumar
相关产品推荐
相关产品推荐

