【发布时间】:2021-03-16 06:38:19
【问题描述】:
我添加了以下代码:
var counters: Map[String, Int] = Map()
val results = rdd.filter(l => l.contains("xyz")).map(l => mapEvent(l)).filter(r => r.isDefined).map (
i => {
val date = i.get.getDateTime.toString.substring(0, 10)
counters = counters.updated(date, counters.getOrElse(date, 0) + 1)
}
)
我想在一次迭代中获取 RDD 中不同日期的计数。但是当我运行它时,我收到消息说:
No implicits found for parameters evidence$6: Encoder[Unit]
所以我添加了这一行:
implicit val myEncoder: Encoder[Unit] = org.apache.spark.sql.Encoders.kryo[Unit]
但是我得到了这个错误。
Exception in thread "main" java.lang.ExceptionInInitializerError
at com.xyz.SparkBatchJob.main(SparkBatchJob.scala)
Caused by: java.lang.UnsupportedOperationException: Primitive types are not supported.
at org.apache.spark.sql.Encoders$.genericSerializer(Encoders.scala:200)
at org.apache.spark.sql.Encoders$.kryo(Encoders.scala:152)
我该如何解决这个问题?或者有没有更好的方法在单次迭代(O(N) 时间)中获得我想要的计数?
【问题讨论】:
-
不要在转换的匿名函数中放入可变变量,只收集第一次转换的结果:i.get.getDateTime.toString.substring(0, 10)
-
@EmiCareOfCell44 - 不知道你的意思。请举例。谢谢。
标签: scala apache-spark