【问题标题】:Why do I get a type mismatch error when using a UDF that returns an object of type Option[Long]?为什么在使用返回 Option[Long] 类型对象的 UDF 时会出现类型不匹配错误?
【发布时间】:2020-09-07 13:59:13
【问题描述】:

我正在尝试在 Scala 中编写一个处理空值的用户定义函数 (UDF)。对于我的示例,如果值不为空,我将尝试返回列的纪元。我发现 Option[] 用于从 udf 返回空值。

这是我的 UDF:

def to_epoch(date: Timestamp) : Option[Long] = {
    if(date != null) {
        Option.apply(date.getTime)
    } else {
        Option.empty
    }
}

val toEpoch: (Timestamp => Option[Long]) => UserDefinedFunction = udf((_: Timestamp => Option[Long]))

我正在从如下读取的文件中创建一个数据框,并且我想添加“dateEpoch”列。我不知道如何让它处理我的udf返回的Option[Long]

spark.read
     .schema(ListeningStatsSchema.schema)
     .json(location)
     .withColumn("dateEpoch", toEpoch(col("EventTS"))

我得到的错误是:

type mismatch;
 found   : org.apache.spark.sql.Column
 required: java.sql.Timestamp => Option[Long]
            .withColumn("opd", toEpoch(col("event_TS")))

【问题讨论】:

  • 你需要使用 Some(date.getTime) 和 None 代替 Option.apply 选项和 None 代替 Option.empty 或 Option.empty[Long],你得到的错误是因为在 with 列中,您没有提供正确的功能。
  • 函数 toEpoch 采用函数 Timestamp => Option[Long] 并且您正在提供列,您应该在那里传递一个采用 Timestamp 并返回 Option[Long] 的函数
  • 问题的标题很可能需要修改@Oli

标签: scala apache-spark user-defined-functions option


【解决方案1】:

您得到的错误意味着您定义的函数需要 Timestamp(参见 REPL 提供的类型)。但是您提供的是Column,因此出现了错误。 Column 类型是您使用 Spark SQL 操作的主要类型。您可以使用预定义的函数和运算符(例如,可以使用 + 添加列)或 UDF,但不能使用常规的 scala 函数。

要修复您的代码,您需要使用 udf 函数将您的函数转换为 spark UDF。你可以这样做:

val to_epoch_udf = udf(to_epoch _)

// And we can try it:
spark.range(1).select(to_epoch_udf(current_timestamp)).show

给出:

+------------------------+
|UDF(current_timestamp())|
+------------------------+
|1599492185730           |
+------------------------+

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2017-02-09
    • 1970-01-01
    • 2020-01-12
    • 1970-01-01
    • 1970-01-01
    • 2017-11-30
    • 2016-03-29
    • 2017-03-14
    相关资源
    最近更新 更多