【问题标题】:Scala: java.lang.UnsupportedOperationException: Primitive types are not supportedScala:java.lang.UnsupportedOperationException:不支持原始类型
【发布时间】: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


【解决方案1】:

Spark RDD 是分布式集合的表示。当您将 map 函数应用于 RDD 时,用于操作集合的函数将在整个集群中执行,因此在 map 函数范围之外创建变量是没有意义的。

在您的代码中,问题在于您没有返回任何值,而是试图改变结构,因此编译器推断转换后新创建的 RDD 是 RDD[Unit]。

如果您需要通过 Spark 操作创建 Map,则必须创建 pairRDD,然后应用 reduce 操作。

包括 rdd 的类型和 mapEvent 函数,看看它是如何完成的。

Spark 使用转换和操作构建 DAG,它不会对数据进行两次处理。

【讨论】:

  • 谢谢。这是有道理的,但是我的问题第二部分的答案是什么:“有没有更好的方法可以在一次迭代中获得计数?”
猜你喜欢
  • 2021-07-06
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2013-03-13
相关资源
最近更新 更多