【问题标题】:Scala Spark dataframe : Task not serilizable exception even with Broadcast variablesScala Spark数据帧:即使使用广播变量,任务也无法序列化异常
【发布时间】:2016-05-09 10:14:34
【问题描述】:

这有效(df:数据框)

val filteredRdd = df.rdd.zipWithIndex.collect { case (r, i) if i >= 10 => r }

这不是

val start=10
val filteredRdd = df.rdd.zipWithIndex.collect { case (r, i) if i >= start => r }

我尝试使用广播变量,但即使这样也没有用

 val start=sc.broadcast(1)
 val filteredRdd = df.rdd.zipWithIndex.collect { case (r, i) if i >= start.value => r }

我收到 Task Not Serializable 异常。任何人都可以解释为什么即使使用广播变量它也会失败。

org.apache.spark.SparkException: Task not serializable
at org.apache.spark.util.ClosureCleaner$.ensureSerializable(ClosureCleaner.scala:304)
at org.apache.spark.util.ClosureCleaner$.org$apache$spark$util$ClosureCleaner$$clean(ClosureCleaner.scala:294)
at org.apache.spark.util.ClosureCleaner$.clean(ClosureCleaner.scala:122)
at org.apache.spark.SparkContext.clean(SparkContext.scala:2055)
at org.apache.spark.rdd.RDD$$anonfun$collect$2.apply(RDD.scala:959)
at org.apache.spark.rdd.RDD$$anonfun$collect$2.apply(RDD.scala:958)
at org.apache.spark.rdd.RDDOperationScope$.withScope(RDDOperationScope.scala:150)
at org.apache.spark.rdd.RDDOperationScope$.withScope(RDDOperationScope.scala:111)
at org.apache.spark.rdd.RDD.withScope(RDD.scala:316)
at org.apache.spark.rdd.RDD.collect(RDD.scala:958)
at $iwC$$iwC$$iwC$$iwC$$iwC$$iwC$$iwC$$iwC$$iwC$$iwC$$iwC$$$$$$fa17825793f04f8d2edd8765c45e2a6c$$$$wC$$iwC$$iwC$$iwC$$iwC$$iwC$$iwC$$iwC$$iwC$$iwC$$iwC$$iwC$$iwC$$iwC$$iwC$$iwC$$iwC$$iwC$$iwC$$iwC$$iwC$$iwC$$iwC$$iwC$$iwC$$iwC$$iwC$$iwC$$iwC$$iwC.(:172)
at $iwC

【问题讨论】:

  • 如果start 是某些(不可序列化)类的一部分,这可能会发生 - 是吗?您能否提供更详细的代码示例?
  • 我在 zipwithIndex 操作之前定义了start。我没有课。这是一个脚本,直接编写用于解释。

标签: scala apache-spark spark-dataframe


【解决方案1】:

您使用的基本结构看起来很可靠。这是一个类似的代码 sn-p 确实 工作。请注意,它使用 broadcast 并使用 map 方法中的广播值 - 类似于您的代码。

scala> val dat = sc.parallelize(List(1,2,3))
dat: org.apache.spark.rdd.RDD[Int] = ParallelCollectionRDD[0] at parallelize at <console>:24
scala> val br = sc.broadcast(10)
br: org.apache.spark.broadcast.Broadcast[Int] = Broadcast(2)

scala> dat.map(br.value * _)
res2: org.apache.spark.rdd.RDD[Long] = MapPartitionsRDD[1] at map at <console>:29

scala> res2.collect
res3: Array[Int] = Array(10, 20, 30)

因此,这可能有助于您验证您的一般方法。

我怀疑您的问题出在脚本中的其他变量上。尝试先在新的 spark-shell 会话中剥离所有内容,然后通过消除过程找出罪魁祸首。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2021-03-04
    • 1970-01-01
    • 1970-01-01
    • 2015-12-16
    • 1970-01-01
    相关资源
    最近更新 更多