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

Spark 1.6 Scala:移除DataFrame数组列中的Null值

在Spark 1.6中过滤数组列的Null元素

针对你这个包含id字符串列和desc结构体数组列的DataFrame,由于Spark 1.6没有内置的array_remove这类高版本函数,我给你提供两种实用的方法来剔除desc数组里的Null元素:

方法一:自定义UDF(最直接的方式)

写一个简单的自定义函数来过滤数组中的Null值,逻辑清晰且易于实现:

import org.apache.spark.sql.functions.udf
import org.apache.spark.sql.types._

// 定义过滤数组Null元素的UDF
val filterNullsInArray = udf((arr: Seq[Row]) => {
  arr.filter(_ != null)
})

// 将UDF应用到desc列,直接替换原列或者新增列都可以
val filteredDF = originalDF.withColumn("desc", filterNullsInArray($"desc"))

补充说明:

因为desc的元素是结构体类型,在Spark 1.6中数组元素会被映射为Row对象,所以UDF的输入参数设为Seq[Row];filter(_ != null)会直接筛掉数组里的Null元素,保留所有非Null的结构体。

方法二:Explode拆分 + GroupBy聚合

如果你需要对数组元素做更复杂的校验,这种拆分再聚合的方式会更灵活:

import org.apache.spark.sql.functions._

// 第一步:把desc数组拆成单独的行,同时保留id关联
val explodedDF = originalDF.select($"id", explode($"desc").alias("desc_element"))

// 第二步:过滤掉Null的结构体行
val filteredExplodedDF = explodedDF.filter($"desc_element".isNotNull)

// 第三步:按id分组,重新组装成过滤后的数组
val filteredDF = filteredExplodedDF.groupBy($"id")
  .agg(collect_list($"desc_element").alias("desc"))

补充说明:

  • explode会将数组的每个元素拆成独立行,包括Null元素;
  • 用filter语句剔除Null行后,再通过collect_list聚合回数组,最终得到不含Null的desc列。

这两种方法都能满足你的需求,你可以根据实际场景选择合适的方案。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.27 03:50:39