【问题标题】:Cannot resolve task not serializable [org.apache.spark.SparkException: Task not serializable] Spark Scala RDD无法解决不可序列化的任务 [org.apache.spark.SparkException: Task not serializable] Spark Scala RDD
【发布时间】:2020-09-22 23:04:49
【问题描述】:

当我尝试创建类的对象并调用特定方法 newRDD 和 blah 时,我不断收到以下错误堆栈跟踪

I create a spark shell by importing the jar and run the following in spark-shell

spark-shell --master=yarn --jars=sample_jar.jar --files database.cfg

scala> val reader = new Sample(spark)
scala> val a = reader.buildFileRDD("/xyz/path")

org.apache.spark.SparkException: Task not serializable
  at org.apache.spark.util.ClosureCleaner$.ensureSerializable(ClosureCleaner.scala:298)
  at org.apache.spark.util.ClosureCleaner$.org$apache$spark$util$ClosureCleaner$$clean(ClosureCleaner.scala:288)
  at org.apache.spark.util.ClosureCleaner$.clean(ClosureCleaner.scala:108)
  at org.apache.spark.SparkContext.clean(SparkContext.scala:2294)
  at org.apache.spark.rdd.RDD$$anonfun$filter$1.apply(RDD.scala:387)
  at org.apache.spark.rdd.RDD$$anonfun$filter$1.apply(RDD.scala:386)
  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:362)
  at org.apache.spark.rdd.RDD.filter(RDD.scala:386)
  at Sample.newRDDscala(Sample.scala:117)
  ... 48 elided
Caused by: java.io.NotSerializableException:

如何解决此错误?

【问题讨论】:

  • 你能发布这个代码 - DatabaseUtils ?? & 你在哪里使用这个 - dbObj
  • @Srinivas 添加了更多详细信息
  • 将 dbObj 更改为 def 或在此之前添加大小写 - DatabaseUtils ?试试
  • @Srinivas 没明白。可以举个例子吗?
  • 经验法则是,如果您使用 scala 类的对象或 spark 闭包内的对象,如 rdd.map/filter/mapPartitions 等,则必须对其进行序列化,希望对您有所帮助:)

标签: scala apache-spark apache-spark-sql rdd


【解决方案1】:

从堆栈跟踪看来,您在闭包内使用DatabaseUtils 的对象,因为DatabaseUtils 不可序列化,因此无法通过n/w 传输,请尝试序列化DatabaseUtils。另外,你可以让DatabaseUtilsscala object

.. DatabaseUtils extends Serializable

【讨论】:

  • 你能扩展你所说的“另外,你可以制作 DatabaseUtils scala 对象”的意思吗?举个例子?
  • 而不是class DatabaseUtils(url: String, username: String, password: String) 使用object DatabaseUtils 并将所有字段移动到实用程序。请注意,这与serializableException无关
  • 我当然可以试一试,但是如何解决由val fs = FileSystem.get(sprk.sparkContext.hadoopConfiguration) 引起的第二个堆栈跟踪?
【解决方案2】:

更改DatabaseUtils 代码如下和内部类示例删除变量dbConfig 和url,添加此val dbObj = new DatabaseUtils(ConfigFactory.parseFile(new File(config)))

class DatabaseUtils(url: String, username: String, password: String) {

  val driver = "com.mysql.jdbc.Driver"


  def executeSelectQuery(qry: String): List[String] = {

    var dbString : ArrayBuffer[String] = ArrayBuffer.empty[String]
    var conn:Connection = null
    try {
      Class.forName(driver)
      conn = DriverManager.getConnection(url, username, password)
      val statement = conn.createStatement
      val rs = statement.executeQuery(qry)

      while (rs.next) dbString += rs.getString("db_string")

    } catch {
      case e: Exception => e.printStackTrace
    }
    finally {
      conn.close()
    }
    dbString.toList
  }
}

object DatabaseUtils {
 def apply(dbConfig:Config): DatabaseUtils = {
    val url =  "jdbc:mysql://" + dbConfig.getString("db.host") +":"+ dbConfig.getString("db.port") + "/" + dbConfig.getString("db.database") + "?useSSL=false"
    new DatabaseUtils(url, dbConfig.getString("db.username") ,dbConfig.getString("db.password"))
  }
}

【讨论】:

  • 您能发布完整的代码吗?您的代码似乎不正确并显示编译错误。
  • 没有代码是正确的。通过 Spark-shell 测试时它工作正常。只有当我手动创建一个 jar 并将 jar 传递给 spark shell 时,它才会出错。
  • 尽量减少类中的变量并将其移动到对象或函数中
猜你喜欢
  • 1970-01-01
  • 2016-12-14
  • 2023-04-02
  • 1970-01-01
  • 1970-01-01
  • 2016-02-07
  • 1970-01-01
  • 1970-01-01
  • 2015-05-31
相关资源
最近更新 更多