【问题标题】:Scala Spark Task not serializable error in code代码中的Scala Spark Task不可序列化错误
【发布时间】:2020-02-20 16:15:08
【问题描述】:

不确定下面的代码有什么问题,但它会抛出

org.apache.spark.SparkException: Task not serializable

错误。谷歌搜索了错误,但无法解决。

以下是代码:(可以通过创建新的 Scala 笔记本复制粘贴并在 community.cloud.databricks.com 上执行)

    import com.google.gson._    
     object TweetUtils {
          case class Tweet(
               id : String,
               user : String,
               userName : String,
               text : String,
               place : String,
               country : String,
               lang : String
          ) 

       def parseFromJson(lines:Iterator[String]):Iterator[Tweet] = {
            val gson = new Gson
            lines.map( line => gson.fromJson(line, classOf[Tweet]))     
       }

       def loadData(): RDD[Tweet] = { 
           val pathToFile = "/FileStore/tables/reduced_tweets-57570.json"
           sc.textFile(pathToFile).mapPartitions(parseFromJson(_))
       }

       def tweetsByUser(): RDD[(String, Iterable[Tweet])] = {
           val tweets = loadData
           tweets.groupBy(_.user)    
       }  
   } 

   val res = TweetUtils.tweetsByUser()
   res.collect().take(5).foreach(println)

以下是详细的错误信息:

at org.apache.spark.util.ClosureCleaner$.ensureSerializable(ClosureCleaner.scala:403)
    at org.apache.spark.util.ClosureCleaner$.org$apache$spark$util$ClosureCleaner$$clean(ClosureCleaner.scala:393)
    at org.apache.spark.util.ClosureCleaner$.clean(ClosureCleaner.scala:162)
    at org.apache.spark.SparkContext.clean(SparkContext.scala:2548)
    at org.apache.spark.rdd.RDD$$anonfun$mapPartitions$1.apply(RDD.scala:827)
    at org.apache.spark.rdd.RDD$$anonfun$mapPartitions$1.apply(RDD.scala:826)
    at org.apache.spark.rdd.RDDOperationScope$.withScope(RDDOperationScope.scala:151)
    at org.apache.spark.rdd.RDDOperationScope$.withScope(RDDOperationScope.scala:112)
    at org.apache.spark.rdd.RDD.withScope(RDD.scala:392)
    at org.apache.spark.rdd.RDD.mapPartitions(RDD.scala:826)
    at line05db6d250e4b42e2b2c1d6b97ba83df533.$read$$iw$$iw$$iw$$iw$$iw$$iw$$iw$$iw$TweetUtils$.loadData(command-3696793732897971:22)
    at line05db6d250e4b42e2b2c1d6b97ba83df533.$read$$iw$$iw$$iw$$iw$$iw$$iw$$iw$$iw$TweetUtils$.tweetsByUser(command-3696793732897971:25)
    at line05db6d250e4b42e2b2c1d6b97ba83df533.$read$$iw$$iw$$iw$$iw$$iw$$iw$$iw$$iw.<init>(command-3696793732897971:30)
    at line05db6d250e4b42e2b2c1d6b97ba83df533.$read$$iw$$iw$$iw$$iw$$iw$$iw$$iw.<init>(command-3696793732897971:84)
    at line05db6d250e4b42e2b2c1d6b97ba83df533.$read$$iw$$iw$$iw$$iw$$iw$$iw.<init>(command-3696793732897971:86)
    at line05db6d250e4b42e2b2c1d6b97ba83df533.$read$$iw$$iw$$iw$$iw$$iw.<init>(command-3696793732897971:88)
    at line05db6d250e4b42e2b2c1d6b97ba83df533.$read$$iw$$iw$$iw$$iw.<init>(command-3696793732897971:90)
    at line05db6d250e4b42e2b2c1d6b97ba83df533.$read$$iw$$iw$$iw.<init>(command-3696793732897971:92)
    at line05db6d250e4b42e2b2c1d6b97ba83df533.$read$$iw$$iw.<init>(command-3696793732897971:94)

提前致谢,

斯里

【问题讨论】:

  • 您能否尝试将case class Tweet 移出 TweetUtils 对象。 TweetUtils 不可序列化。
  • 您好 Artem Aliev,感谢您的回复。我尝试将“案例类推文”移到“对象”之外,但我面临同样的问题“任务不可序列化”。

标签: scala apache-spark


【解决方案1】:

最后,有效的方法是将“Artem Aliev”和“Partha”的建议一起实施。即通过将“案例类 Tweet”移到“TweetUtils 对象”之外并扩展对象“对象 TweetUtils 扩展可序列化”

谢谢你们俩。

【讨论】:

    【解决方案2】:

    将您的TweetUtils 对象设为Serializable,它应该可以工作:

    object TweetUtils extends Serializable

    【讨论】:

    • 嗨,Partha,我尝试使对象可序列化。但我收到错误“格式错误的类名”
    猜你喜欢
    • 2017-09-21
    • 1970-01-01
    • 2020-02-04
    • 2017-11-19
    • 2020-09-22
    • 1970-01-01
    • 2015-12-16
    • 1970-01-01
    • 2021-08-12
    相关资源
    最近更新 更多