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

