无需Pandas从Spark DataFrame中移除空值的实现方案
Spark原生解决方案:移除Struct中Null值字段
问题说明
现有Spark DataFrame处理逻辑中,COL1、COL2字段可能存在Null值,需要生成的PROP字段仅保留非Null的子字段;当两个字段均为Null时,PROP输出为空对象。此前尝试用map_filter处理Struct类型的PROP未生效,使用Pandas处理会因内存限制溢出,需Spark原生解决方案。
原始处理代码
from pyspark.sql import functions as F new_DF = old_DF\ .select(col("id"), col("COL1"), col("COL2"), ).distinct() new_JSON_DF = new_DF\ .withColumn("PROP", struct(col("COL1"), col("COL2")))\ .drop("COL1", "COL2")
样例输入数据
id COL1 COL2 1 null def 2 abc null 3 null null
期望输出JSON格式
{ "id": 1, "PROP": { "COL2": "def" } }, { "id": 2, "PROP": { "COL1": "abc" } }, { "id": 3, "PROP": {} }
解决方案
核心问题说明
之前使用map_filter无效的原因是:struct()生成的是StructType数据,而map_filter仅支持处理MapType数据,类型不匹配导致过滤逻辑未执行。
通用实现方案(适用于任意数量字段)
通过create_map将字段转为MapType,再用map_filter过滤Null值,最终输出的Map转JSON后与目标格式一致:
from pyspark.sql import functions as F # 原始去重逻辑保留 new_DF = old_DF.select(F.col("id"), F.col("COL1"), F.col("COL2")).distinct() # 构建PROP:先转Map,过滤Null值 new_JSON_DF = new_DF.withColumn( "PROP", F.map_filter( # 将COL1、COL2转为键值对Map F.create_map( F.lit("COL1"), F.col("COL1"), F.lit("COL2"), F.col("COL2") ), # 过滤值不为Null的条目 lambda k, v: v.isNotNull() ) ).drop("COL1", "COL2")
字段较少时的简化方案
如果仅处理少量字段,可通过when条件动态构建Struct:
from pyspark.sql import functions as F new_DF = old_DF.select(F.col("id"), F.col("COL1"), F.col("COL2")).distinct() new_JSON_DF = new_DF.withColumn( "PROP", F.when( F.col("COL1").isNotNull() & F.col("COL2").isNotNull(), F.struct(F.col("COL1"), F.col("COL2")) ).when( F.col("COL1").isNotNull(), F.struct(F.col("COL1")) ).when( F.col("COL2").isNotNull(), F.struct(F.col("COL2")) ).otherwise( F.struct() # 空Struct对应JSON空对象 ) ).drop("COL1", "COL2")
验证结果
两种方案处理样例数据后,输出的JSON格式均与期望一致,且完全基于Spark原生API,无内存溢出风险。
内容的提问来源于stack exchange,提问作者dontgimmehope
相关产品推荐
相关产品推荐

