【问题标题】:java.io.NotSerializableException: org.apache.spark.sql.Column when creating a new column conditionally, using a map of UDFsjava.io.NotSerializableException: org.apache.spark.sql.Column 使用 UDF 映射有条件地创建新列时
【发布时间】:2018-08-23 21:05:05
【问题描述】:

我有一个带有 startTime 和一些特征向量的设备 ID 数据,需要基于 hourweekday_hour 进行合并。样本数据如下:

+-----+-------------------+--------------------+
|hh_id|          startTime|                hash|
+-----+-------------------+--------------------+
|dev01|2016-10-10 00:01:04|(1048576,[121964,...|
|dev02|2016-10-10 00:17:45|(1048576,[121964,...|
|dev01|2016-10-10 00:18:01|(1048576,[121964,...|
|dev10|2016-10-10 00:19:48|(1048576,[121964,...|
|dev05|2016-10-10 00:20:00|(1048576,[121964,...|
|dev08|2016-10-10 00:45:13|(1048576,[121964,...|
|dev05|2016-10-10 00:56:25|(1048576,[121964,...|

这些特征基本上是 SparseVectors,由自定义函数合并而成。当我尝试通过以下方式创建 key 列时:

val columnMap = Map("hour" -> hour($"startTime"), "weekday_hour" -> getWeekdayHourUDF($"startTime"))
val grouping = "hour"
val newDF = oldDF.withColumn("dt_key", columnMap(grouping))

我收到了java.io.NotSerializableException。完整的堆栈跟踪如下:

Caused by: java.io.NotSerializableException: org.apache.spark.sql.Column
Serialization stack:
    - object not serializable (class: org.apache.spark.sql.Column, value: hour(startTime))
    - field (class: scala.collection.immutable.Map$Map3, name: value1, type: class java.lang.Object)
    - object (class scala.collection.immutable.Map$Map3, Map(hour -> hour(startTime), weekday_hour -> UDF(startTime), none -> 0))
    - field (class: linef03f4aaf3a1c4f109fce271f7b5b1e30104.$read$$iw$$iw$$iw$$iw$$iw$$iw$$iw$$iw$$iw$$iw, name: groupingColumnMap, type: interface scala.collection.immutable.Map)
    - object (class linef03f4aaf3a1c4f109fce271f7b5b1e30104.$read$$iw$$iw$$iw$$iw$$iw$$iw$$iw$$iw$$iw$$iw, linef03f4aaf3a1c4f109fce271f7b5b1e30104.$read$$iw$$iw$$iw$$iw$$iw$$iw$$iw$$iw$$iw$$iw@4f1f9a63)
    - field (class: linef03f4aaf3a1c4f109fce271f7b5b1e30104.$read$$iw$$iw$$iw$$iw$$iw$$iw$$iw$$iw$$iw, name: $iw, type: class linef03f4aaf3a1c4f109fce271f7b5b1e30104.$read$$iw$$iw$$iw$$iw$$iw$$iw$$iw$$iw$$iw$$iw)
    - object (class linef03f4aaf3a1c4f109fce271f7b5b1e30104.$read$$iw$$iw$$iw$$iw$$iw$$iw$$iw$$iw$$iw, linef03f4aaf3a1c4f109fce271f7b5b1e30104.$read$$iw$$iw$$iw$$iw$$iw$$iw$$iw$$iw$$iw@207d6d1e)

但是当我尝试在不显式创建列的情况下使用 if-else 执行相同的逻辑时,我不会遇到任何此类错误。

val newDF = if(groupingKey == "hour") {
  oldDF.withColumn("dt_key", hour($"startTime")
} else {
  oldDF.withColumn("dt_key", getWeekdayHourUDF($"startTime")
}

使用 Map-way 会非常方便,因为可能有更多类型的密钥提取方法。请帮助我找出导致此问题的原因。

【问题讨论】:

  • 你应该写一个 udf 函数来创建地图
  • @RameshMaharjan 你的意思是,我需要创建一个包含要应用的函数的映射的 UDF,而不是创建 UDF 的映射?
  • 我无法在 Spark 1.6 或 Spark 2.2 中重现这一点。你确定UDF没有问题吗? columnMap 的类型是什么?应该是Map[String, Column]
  • @philantrovert 是的,columnMap 的类型是scala.collection.immutable.Map[String,org.apache.spark.sql.Column]。我认为 UDF 没有问题,但这里是它的代码。 val localDateTime = ts.toLocalDateTime 然后(localDateTime.getDayOfWeek.getValue - 1)*24 + localDateTime.getHour val getWeekdayHourUDF = udf(getWeekdayHour _)

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


【解决方案1】:

内置功能时

您可以通过使用when 内置函数来实现您的要求

val groupingKey = //"hour" or "weekday_hour"
import org.apache.spark.sql.functions._
df.withColumn("dt_key", 
     when(lit(groupingKey) === "hour", hour($"startTime"))
     .when(lit(groupingKey) === "weekday_hour", getWeekdayHourUDF($"startTime"))
     .otherwise(lit(0)))).show(false)

udf 函数

或者,您可以创建一个udf 函数以创建地图列

import org.apache.spark.sql.functions._
def mapUdf = udf((hour: Int, weekdayhour: Int, groupingKey: String) => 
      if(groupByKey.equalsIgnoreCase("hour")) hour 
      else if(groupByKey.equalsIgnoreCase("weekday_hour")) weekdayhour 
      else 0)

并将其用作

val newDF = oldDF.withColumn("dt_key",
                  mapUdf(hour($"startTime"), 
                         getWeekdayHourUDF($"startTime"),
                         lit(groupingKey)))

希望回答对你有帮助

【讨论】:

  • 我认为这个问题存在误解。数据框有startTime 列,dt_key 列是根据groupingKey 的值导出的,其值可以是hourweekday_hour。最后,我不需要列映射,我只需要一个列作为结果。
  • 那么groupingKey是什么?是另一列吗?
  • 不,它只是一个字符串变量,其值为"hour""weekday_hour"。它定义了对整个数据进行分组的类型。我指的另一个函数hourspark.sql.functions.hour,它从时间戳中提取小时。
  • 啊!我知道了。您已将 if-else 转移到 mapUDF 中。但问题是,我无法在这个方法中使用 Spark 提供的原生函数。 :( 我希望有其他更简单的方法可以做到这一点。我尝试使用transient 关键字作为columnMap,但即使这样也不起作用。
  • 我又更新了答案 :) 希望这次我如你所愿地回答了
【解决方案2】:

可能有点晚了,但我使用的是 Spark 2.4.6,无法重现该问题。我猜代码调用columnMap 用于多个键。如果您提供一个易于重现的示例,包括数据(1 行数据集就足够了),它会有所帮助。但是,正如堆栈跟踪所说,Column 类确实不是Serializable,我将根据我目前的理解尝试详细说明。

TLDR;一种简单的规避方法是将vals 转换为defs。


我相信为什么用when case 或 UDF 表达同样的事情已经很清楚了。

第一次尝试:这样的事情可能行不通的原因是(a)Column 类不可序列化(我认为这是一个有意识的设计选择,因为它的预期角色是Spark API),并且(b)表达式中没有任何内容

oldDF.withColumn("dt_key", columnMap(grouping))

这告诉 Spark 对于withColumn 的第二个参数实际具体的Column 是什么,这意味着具体的Map[String, Column] 对象将需要通过网络发送给执行程序,当这样的异常会被抚养。

第二次尝试:第二次尝试有效的原因是,定义DataFrame 所需的有关groupingKey 参数的相同决定可能完全发生在驱动程序上。


考虑使用 DataFrame API 作为查询构建器的 Spark 代码,或者包含执行计划的东西,而不是数据本身,这会有所帮助。一旦您对其调用操作(writeshowcount 等),Spark 就会生成将任务发送给执行程序的代码。此时,实现DataFrame/Dataset 所需的所有信息必须已经在查询计划中正确编码,或者需要可序列化以便可以通过网络发送。

def 通常会解决这类问题,因为

def columnMap: Map[String, Column] = Map("a" -> hour($"startTime"), "weekday_hour" -> UDF($"startTime"))

不是具体的Map 对象本身,而是每次调用它时创建一个新的Map[String, Column] 的东西,比如在每个碰巧执行涉及此Map 的任务的执行者处.

Thisthis 似乎是该主题的好资源。我承认我明白为什么要使用 Function 之类的

val columnMap = () => Map("a" -> hour($"startTime"), "b" -> UDF($"startTime"))

然后columnMap()("a") 会起作用,因为反编译的字节码显示scala.Functions 被定义为Serializable 的具体实例,但我不明白为什么def 起作用,因为那看起来不对他们来说就是这样。无论如何,我希望这会有所帮助。

【讨论】:

  • 谢谢!将val 更改为def 为我做了这件事!
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 2021-12-11
  • 1970-01-01
  • 1970-01-01
  • 2020-08-13
  • 2018-01-10
  • 1970-01-01
  • 2021-01-04
相关资源
最近更新 更多