这里的问题与 Job 中使用的 avro.Schema 类的不可序列化有关。当您尝试从 map 函数内的代码中引用架构对象时,将引发异常。
例如,如果您尝试执行以下操作,您将收到“Task not serializable”异常:
val schema = new Schema.Parser().parse(new File(jsonSchema))
rdd.map(t => {
// reference to the schema object declared outside
val record = new GenericData.Record(schema)
val schema = new Schema.Parser().parse(new File(jsonSchema))
// The schema above should not be used in closures, it's for other purposes
rdd.map(t => {
// create a new Schema object
val innserSchema = new Schema.Parser().parse(new File(jsonSchema))
val record = new GenericData.Record(innserSchema)
由于您不希望为您处理的每条记录解析 avro 架构,因此更好的解决方案是在分区级别解析架构。以下也有效:
val schema = new Schema.Parser().parse(new File(jsonSchema))
// The schema above should not be used in closures, it's for other purposes
rdd.mapPartitions(tuples => {
// create a new Schema object
val innserSchema = new Schema.Parser().parse(new File(jsonSchema))
tuples.map(t => {
val record = new GenericData.Record(innserSchema)
// this closure will be bundled together with the outer one
// (no serialization issues)
只要您提供对 jsonSchema 文件的可移植引用,上面的代码就可以工作,因为 map 函数将由多个远程执行程序执行。它可以是对 HDFS 中文件的引用,也可以与 JAR 中的应用程序一起打包(在后一种情况下,您将使用类加载器函数来获取其内容)。
对于那些尝试将 Avro 与 Spark 一起使用的人,请注意仍然存在一些未解决的编译问题,您必须在 Maven POM 上使用以下导入: