【问题标题】:Task Not Serializable exception when trying to write a rdd of type Generic Record尝试编写 Generic Record 类型的 rdd 时,Task Not Serializable 异常
【发布时间】:2017-11-15 10:17:11
【问题描述】:
val file = File.createTempFile("temp", ".avro")
val schema = new Schema.Parser().parse(st)
val datumWriter = new GenericDatumWriter[GenericData.Record](schema)
val dataFileWriter = new DataFileWriter[GenericData.Record](datumWriter)
dataFileWriter.create(schema , file)
rdd.foreach(r => {
  dataFileWriter.append(r)
})
dataFileWriter.close()

我有一个GenericData.Record 类型的DStream,我正在尝试以Avro 格式写入HDFS,但我收到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:2062)
at org.apache.spark.rdd.RDD$$anonfun$foreach$1.apply(RDD.scala:911)
at org.apache.spark.rdd.RDD$$anonfun$foreach$1.apply(RDD.scala:910)
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.foreach(RDD.scala:910)
at KafkaCo$$anonfun$main$3.apply(KafkaCo.scala:217)
at KafkaCo$$anonfun$main$3.apply(KafkaCo.scala:210)
at org.apache.spark.streaming.dstream.DStream$$anonfun$foreachRDD$1$$anonfun$apply$mcV$sp$3.apply(DStream.scala:661)
at org.apache.spark.streaming.dstream.DStream$$anonfun$foreachRDD$1$$anonfun$apply$mcV$sp$3.apply(DStream.scala:661)
at org.apache.spark.streaming.dstream.ForEachDStream$$anonfun$1$$anonfun$apply$mcV$sp$1.apply$mcV$sp(ForEachDStream.scala:50)
at org.apache.spark.streaming.dstream.ForEachDStream$$anonfun$1$$anonfun$apply$mcV$sp$1.apply(ForEachDStream.scala:50)
at org.apache.spark.streaming.dstream.ForEachDStream$$anonfun$1$$anonfun$apply$mcV$sp$1.apply(ForEachDStream.scala:50)
at org.apache.spark.streaming.dstream.DStream.createRDDWithLocalProperties(DStream.scala:426)
at org.apache.spark.streaming.dstream.ForEachDStream$$anonfun$1.apply$mcV$sp(ForEachDStream.scala:49)
at org.apache.spark.streaming.dstream.ForEachDStream$$anonfun$1.apply(ForEachDStream.scala:49)
at org.apache.spark.streaming.dstream.ForEachDStream$$anonfun$1.apply(ForEachDStream.scala:49)
at scala.util.Try$.apply(Try.scala:161)
at org.apache.spark.streaming.scheduler.Job.run(Job.scala:39)
at org.apache.spark.streaming.scheduler.JobScheduler$JobHandler$$anonfun$run$1.apply$mcV$sp(JobScheduler.scala:224)
at org.apache.spark.streaming.scheduler.JobScheduler$JobHandler$$anonfun$run$1.apply(JobScheduler.scala:224)
at org.apache.spark.streaming.scheduler.JobScheduler$JobHandler$$anonfun$run$1.apply(JobScheduler.scala:224)
at scala.util.DynamicVariable.withValue(DynamicVariable.scala:57)
at org.apache.spark.streaming.scheduler.JobScheduler$JobHandler.run(JobScheduler.scala:223)
at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1142)
at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:617)
at java.lang.Thread.run(Thread.java:745)
Caused by: java.io.NotSerializableException: org.apache.avro.file.DataFileWriter
Serialization stack:
- object not serializable (class: org.apache.avro.file.DataFileWriter, value: org.apache.avro.file.DataFileWriter@78f132d9)
- field (class: KafkaCo$$anonfun$main$3$$anonfun$apply$1, name: dataFileWriter$1, type: class org.apache.avro.file.DataFileWriter)
- object (class KafkaCo$$anonfun$main$3$$anonfun$apply$1, <function1>)
at org.apache.spark.serializer.SerializationDebugger$.improveException(SerializationDebugger.scala:40)
at org.apache.spark.serializer.JavaSerializationStream.writeObject(JavaSerializer.scala:47)
at org.apache.spark.serializer.JavaSerializerInstance.serialize(JavaSerializer.scala:101)
at org.apache.spark.util.ClosureCleaner$.ensureSerializable(ClosureCleaner.scala:301)

【问题讨论】:

  • 你想通过将 RDD 对象写入 Avro 文件来达到什么目的?您应该在 github.com/databricks/spark-avro 上抢购一下,它可以让您使用 df.write.avro("/tmp/output") 之类的东西将 DataFrame 直接保存到 Avro 格式的文件中

标签: scala apache-spark spark-streaming avro


【解决方案1】:

这里的重点是DataFileWriter是本地资源(绑定到本地文件),所以序列化它没有意义。

调整代码以执行mapPartitions 之类的操作也无济于事,因为这种执行器绑定的方法会将文件写入执行器的本地文件系统。

我们需要使用支持 Spark 分布式特性的实现,例如https://github.com/databricks/spark-avro

使用该库:

给定一些由case class 表示的架构,我们会这样做:

val structuredRDD = rdd.map(record => recordToSchema(record))
val df = structuredRDD.toDF()
df.write.avro(hdfs_path)

【讨论】:

  • 什么是recordToSchema
  • 您编写的用于将记录格式转换为案例类的函数。
  • 您可以提供任何示例,因为我之前没有这样做过
  • @JSR29 好吧,我看到 GenericData.Record 已经是 AVRO 格式了。也许有更聪明的方法可以直接将其保存为 AVRO。 - 我不知道。
  • 这也适用于流媒体吗?问题提到DStreams.
【解决方案2】:

由于 lambdas 必须分布在集群周围才能运行,它们必须只引用可序列化的数据,以便它们可以被序列化、运送到不同的执行器进行部署并作为任务在那里执行。

你可能会做的是:

  • 创建一个新文件并获取它的句柄
  • 使用mapPartitions(而不是map)方法并为每个分区创建一个新的写入器
  • 将文件句柄与您为每个分区创建的编写器一起使用,将分区内的每条消息附加到该文件中
  • 确保在流完全使用时关闭文件句柄

【讨论】:

  • 我在哪里使用 map ,伪代码会很有帮助
  • 虽然一般来说是有效的建议,但我认为mapPartitions 和本地对象在这种情况下不会有帮助。注意dataFileWriter 是如何绑定到本地文件的:dataFileWriter.create(schema , file),使用mapPartitions 我们将在每个执行程序上创建一个本地文件。
  • 这就是为什么我建议传递一个文件处理程序以在流完全消耗而不是当前操作模式时关闭。今天我希望有时间用一些代码来编辑我的答案。
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2018-07-08
  • 2020-11-13
  • 2020-11-08
  • 2018-01-04
  • 1970-01-01
  • 2016-12-14
相关资源
最近更新 更多