如何用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
相关产品推荐
相关产品推荐

