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

如何用PySpark通用代码对比两个DataFrame并提取含空值的行

通用PySpark实现:提取两DataFrame中对应列存在空值的行

问题背景

给定两个结构相同的PySpark DataFrame:

df1

columnA columnB columnC columnD
value1  value7  value13 value20
value2  value8  value14 value21
value3  value9  value15 value22
value4  value10 value16 value23
value5  value11 value17 value24
value6  null    null    value25

df2

columnA columnB columnC columnD
value1  value7  value13 value20
value2  null    value14 value21
null    value9  value15 value22
value4  value10 value16 value23
value5  value11 value17 value24
value6  value12 value18 value25

需求

提取所有满足以下任一条件的行:

  • 该行来自df1,且df2的匹配行中至少有一列是空值(同时df1该行对应列不为空)
  • 该行来自df2,且df1的匹配行中至少有一列是空值(同时df2该行对应列不为空)

最终预期结果:

outputDF

columnA columnB columnC columnD
value2  value8  value14 value21
value3  value9  value15 value22
value6  value12 value18 value25

通用PySpark代码实现

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

# 初始化SparkSession(若未初始化)
spark = SparkSession.builder.appName("CompareDFNulls").getOrCreate()

def extract_mismatch_null_rows(df1, df2):
    # 获取DataFrame的所有列名
    cols = df1.columns
    
    # 构建连接条件:处理空值,确保内容匹配的行能关联
    join_conds = []
    for c in cols:
        # 若某列一方为空,用另一方非空值匹配;双方都为空时视为匹配
        join_conds.append(coalesce(df1[c], df2[c]) == coalesce(df2[c], df1[c]))
    
    # 全外连接两个DataFrame
    joined_df = df1.alias("df1").join(df2.alias("df2"), on=join_conds, how="fullouter")
    
    # 构建筛选条件
    filter_conds = []
    for c in cols:
        # df1的行符合条件:df1该列非空,df2对应列空
        filter_conds.append(col(f"df1.{c}").isNotNull() & col(f"df2.{c}").isNull())
        # df2的行符合条件:df2该列非空,df1对应列空
        filter_conds.append(col(f"df2.{c}").isNotNull() & col(f"df1.{c}").isNull())
    
    # 合并所有筛选条件(满足任一即可)
    final_filter = filter_conds[0]
    for cond in filter_conds[1:]:
        final_filter = final_filter | cond
    
    # 分别提取df1和df2中符合条件的行,过滤全空行
    df1_result = joined_df.filter(final_filter).select([col(f"df1.{c}").alias(c) for c in cols]).dropna(how="all")
    df2_result = joined_df.filter(final_filter).select([col(f"df2.{c}").alias(c) for c in cols]).dropna(how="all")
    
    # 合并结果并去重
    output_df = df1_result.union(df2_result).dropDuplicates()
    
    return output_df

# 调用函数(假设df1、df2已提前创建)
output_df = extract_mismatch_null_rows(df1, df2)
# 查看结果
output_df.show()

代码说明

  1. 连接逻辑:通过coalesce函数处理空值,确保两个DataFrame中内容匹配的行(无论哪一方存在空值)能正确关联。
  2. 筛选逻辑:遍历所有列,检查当前DataFrame列非空但对方对应列空的情况,覆盖所有可能的不匹配空值场景。
  3. 通用性:自动读取DataFrame的列名,无需硬编码,适配任意结构相同的DataFrame。
  4. 结果处理:合并df1和df2的符合条件行并去重,避免重复数据。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.04 05:10:18