【发布时间】:2018-08-23 21:05:05
【问题描述】:
我有一个带有 startTime 和一些特征向量的设备 ID 数据,需要基于 hour 或 weekday_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.getHourval getWeekdayHourUDF = udf(getWeekdayHour _)
标签: scala apache-spark apache-spark-sql user-defined-functions