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

如何用PySpark比较两个DataFrame并提取含空值的行

用PySpark提取两个DataFrame中存在Null值的对应有效行

现有两个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

需求:比较两个DataFrame,提取任一DataFrame中存在Null值的行对应的另一个DataFrame中的有效行,最终结果如下:

outputDF:

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

实现方案

方案一:基于关联匹配的方式

核心逻辑是分别找出两个DataFrame中含Null的行,通过主键关联到另一个DataFrame获取对应有效行,最后合并结果。

from pyspark.sql import SparkSession
from pyspark.sql.functions import col, isnull, array

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

# 构建示例DataFrame
data1 = [
    ("value1", "value7", "value13", "value20"),
    ("value2", "value8", "value14", "value21"),
    ("value3", "value9", "value15", "value22"),
    ("value4", "value10", "value16", "value23"),
    ("value5", "value11", "value17", "value24"),
    ("value6", None, None, "value25")
]
df1 = spark.createDataFrame(data1, ["columnA", "columnB", "columnC", "columnD"])

data2 = [
    ("value1", "value7", "value13", "value20"),
    ("value2", None, "value14", "value21"),
    (None, "value9", "value15", "value22"),
    ("value4", "value10", "value16", "value23"),
    ("value5", "value11", "value17", "value24"),
    ("value6", "value12", "value18", "value25")
]
df2 = spark.createDataFrame(data2, ["columnA", "columnB", "columnC", "columnD"])

# 1. 找出df1中含Null的行,关联df2获取对应有效行
df1_null_keys = df1.filter(array(*[isnull(col(c)) for c in df1.columns])).select("columnD")
df2_valid_rows = df2.join(df1_null_keys, on="columnD", how="inner")

# 2. 找出df2中含Null的行,关联df1获取对应有效行
df2_null_keys = df2.filter(array(*[isnull(col(c)) for c in df2.columns])).select("columnD")
df1_valid_rows = df1.join(df2_null_keys, on="columnD", how="inner")

# 3. 合并结果并去重
outputDF = df1_valid_rows.union(df2_valid_rows).dropDuplicates()

# 查看结果
outputDF.show()

方案二:基于分组聚合的通用方式

如果主键不明确,可通过合并两个DataFrame后分组,直接提取每组中无Null的行。

from pyspark.sql import SparkSession
from pyspark.sql.functions import col, isnull, array, lit, when, max, struct

spark = SparkSession.builder.appName("ExtractNonNullRows").getOrCreate()

# 构建示例DataFrame(同方案一,此处省略)

# 给两个DataFrame添加来源标识
df1_tagged = df1.withColumn("source", lit("df1"))
df2_tagged = df2.withColumn("source", lit("df2"))

# 合并DataFrame
combined_df = df1_tagged.union(df2_tagged)

# 判断行是否包含Null
has_null = array(*[isnull(col(c)) for c in combined_df.columns if c != "source"]).cast("boolean")

# 分组后提取无Null的行
outputDF = combined_df.groupBy("columnD").agg(
    max(when(~has_null, struct(*[c for c in combined_df.columns if c != "source"]))).alias("valid_row")
).select("valid_row.*").filter("valid_row is not null")

outputDF.show()

说明

  • 方案一中使用columnD作为关联主键,实际场景可根据业务调整为多列组合(如on=["columnB", "columnC", "columnD"]);
  • 方案二更通用,无需指定主键,自动提取每组中无Null的有效行,若同一组两行都含Null则会被过滤。

内容的提问来源于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 23:20:41