如何在Scala DataFrame的数组列中用默认值填充空值?
替换Spark DataFrame数组列中的空值
原始数据情况
运行你提供的代码后,DataFrame的结构和数据如下:
import spark.implicits._ val columns=Array("id", "subject") val df1=sc.parallelize(Seq( (1, Array("eng","math",null,null)) )).toDF(columns: _*) df1.printSchema // root // |-- id: integer (nullable = false) // |-- subject: array (nullable = true) // | |-- element: string (containsNull = true) df1.show(false) // +---+------------------------+ // |id |subject | // +---+------------------------+ // |1 |[eng, math, null, null] | // +---+------------------------+
解决方案
下面提供两种可行的替换方式,适配不同Spark版本:
方法1:使用Spark内置transform函数(Spark 2.4+)
利用Spark的数组遍历函数transform,结合coalesce快速替换空值:
import org.apache.spark.sql.functions.{transform, coalesce, lit} // 将数组中的null替换为"unknown",生成新列subject_clean val df_clean = df1.withColumn( "subject_clean", transform($"subject", elem => coalesce(elem, lit("unknown"))) ) df_clean.show(false)
输出结果:
+---+------------------------+------------------------------+ |id |subject |subject_clean | +---+------------------------+------------------------------+ |1 |[eng, math, null, null] |[eng, math, unknown, unknown] | +---+------------------------+------------------------------+
方法2:自定义UDF(兼容Spark 2.4以下版本)
如果你的Spark版本较低,无法使用transform,可以通过自定义UDF处理数组:
import org.apache.spark.sql.functions.udf // 定义UDF:遍历数组,替换null为指定值 val replaceArrayNulls = udf((arr: Array[String]) => arr.map(elem => if (elem == null) "unknown" else elem)) val df_clean = df1.withColumn("subject_clean", replaceArrayNulls($"subject")) df_clean.show(false)
该方法的输出结果和方法1一致。
内容的提问来源于stack exchange,提问作者Shankar Panda
相关产品推荐
相关产品推荐

