【问题标题】:The Spark UDF is not changing the column value from null to 0Spark UDF 没有将列值从 null 更改为 0
【发布时间】:2019-10-03 00:56:45
【问题描述】:

尝试使用下面的 UDF 将 Dataframe 中的 null 替换为 0。 在我可能出错的地方,代码看起来很简单,但没有按预期工作。

我尝试创建一个 UDF,用于替换任何值为 null 的列中的 0。

提前谢谢大家。

//imports

object PlayGround {
def missingValType2(n: Int):Int = {
    if(n == null){
      0
    }else{
      n
    }
  }

   def main(args: Array[String]): Unit = {

    Logger.getLogger("org").setLevel(Level.ERROR)
    val spark = SparkSession
      .builder()
      .appName("PlayGround")
      .config("spark.sql.warehouse.dir", "file:///C:/temp")
      .master("local[*]")
      .getOrCreate()

    val missingValUDFType2 = udf[Int, Int](missingValType2)

     val schema = List(
      StructField("name", types.StringType, false),
      StructField("age", types.IntegerType, true)
    )

    val data = Seq(
      Row("miguel", null),
      Row("luisa", 21)
    )
    val df = spark.createDataFrame(
      spark.sparkContext.parallelize(data),
      StructType(schema)
    )
    df.show(false)
    df.withColumn("ageNullReplace",missingValUDFType2($"age")).show()

  }
}

/**
  * +------+----+
  * |name  |age |
  * +------+----+
  * |miguel|null|
  * |luisa |21  |
  * +------+----+
  *
  * Below is the current output.
  * +------+----+--------------+
  * |  name| age|ageNullReplace|
  * +------+----+--------------+
  * |miguel|null|          null|
  * | luisa|  21|            21|
  * +------+----+--------------+*/

预期输出:

 * +------+----+--------------+
  * |  name| age|ageNullReplace|
  * +------+----+--------------+
  * |miguel|null|             0|
  * | luisa|  21|            21|
  * +------+----+--------------+

【问题讨论】:

  • 您为什么要尝试使用 UDF?你可以在你的 withColumn 中使用 when 来做同样的事情。如果您有可以执行相同操作的本机函数,则不建议使用 UDF
  • 您好,@user2315840 感谢您的回复,是的,我可以做到的。我读了这个mungingdata.com/apache-spark/dealing-with-null 部分用户定义的函数,我认为它是不可取的。让我快速尝试一下您的解决方案。

标签: scala apache-spark apache-spark-sql spark-streaming


【解决方案1】:

不需要UDF。您可以将na.fill 应用于DataFrame 中特定类型列的列表,如下所示:

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

val df = Seq(
  ("miguel", None), ("luisa", Some(21))
).toDF("name", "age")

df.na.fill(0, Seq("age")).show
// +------+---+
// |  name|age|
// +------+---+
// |miguel|  0|
// | luisa| 21|
// +------+---+

【讨论】:

  • 你好 Leo,谢谢你。 None 是否推断为空?我想将值为 null 的年龄替换为 0。或某个随机数。
  • 另外,我在这里学到了一个教训:stackoverflow.com/questions/15777745/… 比较 Int 和 Null 类型的值,使用 '!=' 或 '==' 将始终在 != 的情况下产生 true 或 false 在== 。 if(n != null){
  • 并且还在 sql 方式中做了类似的事情:spark.sql("select *, case when age is null then '0' else age end as extraColumn from df2").show()
  • Spark 中 SQL 的 case when/then/else/end 等价于 when/otherwise,这是另一个答案所暗示的。另一种选择是使用coalesce,例如df.withColumn("age", coalesce($"age", lit(0)))。尽管如此,提议的na.fill 方法提供了将null 替换同时应用于多个列的灵活性。
  • 是 Leo,但是当我使用上述解决方案时,它肯定会起作用,但对我来说 df.withColumn("ageNullReplace", when(col("age").isNull,functions.lit( 0)).otherwise(col("age"))) - 出错时无法解析符号。
【解决方案2】:

您可以将 WithColumn 与下面的 when 条件一起使用 代码未经测试

df.withColumn("ageNullReplace", when(col("age").isNull,lit(0)).otherwise(col(age)))

在上面的代码中,否则不需要仅供参考

希望对你有帮助

【讨论】:

  • 嘿 user2315840,我尝试了完全相同的方法: df.withColumn("ageNullReplace", when(col("age").isNull,functions.lit(0)).otherwise(col("age "))) 但是关键字'when'无法解析。
  • 能不能导入org.apache.spark.sql.functions._
猜你喜欢
  • 2015-08-16
  • 1970-01-01
  • 2019-03-11
  • 2019-08-23
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2022-08-18
  • 1970-01-01
相关资源
最近更新 更多