Scala Spark导出含复杂UNION类型Avro文件报错求助
Spark导出含复杂Union类型的Avro文件解决方案
问题背景
需要将Spark RDD/DataSet导出为包含多类型Union的Avro文件,给定的Avro Schema中identifier字段是int、personal_id_code(record)、person_name(record)三者的Union类型。尝试两种方式均报错:
- 使用Spark Avro Writer时,报错提示仅支持
null+非null的Union类型,不支持多类型Union; - 使用Hadoop Avro API时,因Scala样例类结构与Avro Union不匹配(样例类将Union的三种类型作为同对象的字段,而Avro Union要求三选一),导致类型不匹配报错。
给定的Avro Schema
{ "name": "company","type": "record", "fields": [{ "name": "identifier", "type": [ { "name": "uid","type": "int" },{ "name": "personal_id_code","type": "record", "fields": [ {"name": "code","type": "string"}, {"name": "year_released","type": "string"} ] },{ "name": "person_name","type": "record", "fields": [ {"name": "name","type": "string"}, {"name": "surname","type": "string"} ] } ] },{ "name": "users", "type": { "type": "array", "items": { "name": "userdata", "type": "record", "fields": [{"name": "uid","type": "int"}, {"name": "name","type": "string"}, {"name": "zip","type": "int"}, {"name": "timestamp","type": "long"}, {"name": "properties","type": "int"} ] } } } ] }
对应Scala样例类
case class personal_id_code(code: String, year_released: String) case class person_name(name: String, surname: String) case class identifier(uid: Int, personal_id_code: personal_id_code, name: person_name) case class company(identifier: identifier, users: List[userdata]) case class userdata(uid: Int, name: String, zip: Int, timestamp: Long, properties: Int)
核心问题
- Spark Avro官方库仅支持**
null与单个非null类型**组成的Union,不支持多类型(≥3种)Union; - 原Scala样例类的
identifier是包含三个字段的复合类型,与Avro要求的三选一Union类型结构不匹配,导致Hadoop API无法识别。
可行解决方案
方案1:手动转换为Avro GenericRecord
通过手动构建Avro的GenericRecord,匹配Union类型的三选一规则,将Scala对象转换为符合Schema的Avro结构。
代码实现
import org.apache.avro.Schema import org.apache.avro.generic.{GenericData, GenericRecord} import org.apache.hadoop.io.NullWritable import org.apache.hadoop.mapreduce.Job import org.apache.avro.mapred.AvroKey import org.apache.avro.mapreduce.AvroKeyOutputFormat import org.apache.commons.io.FileUtils import java.io.File // 定义转换函数:将Scala company对象转为Avro GenericRecord def convertToGenericRecord(company: company, schema: Schema): GenericRecord = { val companyRecord = new GenericData.Record(schema) // 处理identifier的Union类型:三选一赋值 val identifierUnionSchema = schema.getField("identifier").schema() val identifierValue = company.identifier match { // 优先判断uid是否有效(根据业务调整非空判断逻辑) case id if id.uid != 0 => id.uid: Integer case id if id.personal_id_code != null => val pidRecord = new GenericData.Record(identifierUnionSchema.getTypes.get(1)) pidRecord.put("code", id.personal_id_code.code) pidRecord.put("year_released", id.personal_id_code.year_released) pidRecord case id if id.name != null => val nameRecord = new GenericData.Record(identifierUnionSchema.getTypes.get(2)) nameRecord.put("name", id.name.name) nameRecord.put("surname", id.name.surname) nameRecord case _ => null } companyRecord.put("identifier", identifierValue) // 处理users数组 val usersItemSchema = schema.getField("users").schema().getElementType val usersRecords = company.users.map { user => val userRecord = new GenericData.Record(usersItemSchema) userRecord.put("uid", user.uid: Integer) userRecord.put("name", user.name) userRecord.put("zip", user.zip: Integer) userRecord.put("timestamp", user.timestamp: java.lang.Long) userRecord.put("properties", user.properties: Integer) userRecord } companyRecord.put("users", new GenericData.Array[GenericRecord](usersRecords.size, usersItemSchema, usersRecords)) companyRecord } // 执行导出 val spark = SparkSession.builder().appName("AvroComplexUnion").master("local[*]").getOrCreate() import spark.implicits._ val avroSchemaStr = """{/* 替换为完整的Avro Schema字符串 */}""" val avroSchema = new Schema.Parser().parse(avroSchemaStr) val outputPath = "C:/Temp/avrooutput" FileUtils.deleteDirectory(new File(outputPath)) val myRdd = List( company(identifier(12345678, null, null), List(userdata(123, "John", 123, 789L, 432), userdata(234, "Paul", 234, 890L, 543))) ).rdd val job = Job.getInstance(spark.sparkContext.hadoopConfiguration) AvroJob.setOutputKeySchema(job, avroSchema) myRdd.map(comp => (new AvroKey(convertToGenericRecord(comp, avroSchema)), NullWritable.get())) .saveAsNewAPIHadoopFile( outputPath, classOf[AvroKey[GenericRecord]], classOf[NullWritable], classOf[AvroKeyOutputFormat[GenericRecord]], job.getConfiguration )
关键点
- 根据业务逻辑判断
identifier的实际类型,选择Union中的对应类型赋值; - 手动构建每个Record和数组,严格匹配Avro Schema的结构。
方案2:使用Avro SpecificRecord生成类
通过Avro工具根据Schema生成Java的SpecificRecord类,利用类型安全的方式转换数据并导出。
步骤1:生成Java类
使用avro-tools命令(或Maven/Gradle插件)根据Avro Schema生成Java类:
# 下载avro-tools.jar后执行 java -jar avro-tools-1.11.0.jar compile schema company.avsc ./src/main/java
生成的类包括Company.java、PersonalIdCode.java、PersonName.java等。
步骤2:Scala中转换并导出
import org.apache.hadoop.io.NullWritable import org.apache.hadoop.mapreduce.Job import org.apache.avro.mapred.AvroKey import org.apache.avro.mapreduce.AvroKeyOutputFormat import org.apache.commons.io.FileUtils import java.io.File // 导入生成的Avro Java类 import com.example.Company import com.example.PersonalIdCode import com.example.PersonName import com.example.Userdata val spark = SparkSession.builder().appName("AvroComplexUnion").master("local[*]").getOrCreate() val outputPath = "C:/Temp/avrooutput" FileUtils.deleteDirectory(new File(outputPath)) val myRdd = List( company(identifier(12345678, null, null), List(userdata(123, "John", 123, 789L, 432), userdata(234, "Paul", 234, 890L, 543))) ).rdd // 转换Scala对象为Avro SpecificRecord val convertedRdd = myRdd.map { comp => val avroCompany = new Company() // 处理identifier Union val id = comp.identifier if (id.uid != 0) { avroCompany.setIdentifier(id.uid: Integer) } else if (id.personal_id_code != null) { val pid = new PersonalIdCode() pid.setCode(id.personal_id_code.code) pid.setYearReleased(id.personal_id_code.year_released) avroCompany.setIdentifier(pid) } else if (id.name != null) { val name = new PersonName() name.setName(id.name.name) name.setSurname(id.name.surname) avroCompany.setIdentifier(name) } // 处理users数组 val avroUsers = comp.users.map { user => val ud = new Userdata() ud.setUid(user.uid: Integer) ud.setName(user.name) ud.setZip(user.zip: Integer) ud.setTimestamp(user.timestamp: java.lang.Long) ud.setProperties(user.properties: Integer) ud }.asJava avroCompany.setUsers(avroUsers) avroCompany } // 导出Avro文件 val job = Job.getInstance(spark.sparkContext.hadoopConfiguration) AvroJob.setOutputKeySchema(job, Company.getClassSchema) convertedRdd.map(avroComp => (new AvroKey(avroComp), NullWritable.get())) .saveAsNewAPIHadoopFile( outputPath, classOf[AvroKey[Company]], classOf[NullWritable], classOf[AvroKeyOutputFormat[Company]], job.getConfiguration )
关键点
- 生成的SpecificRecord类自带类型匹配逻辑,无需手动构建GenericRecord;
- 类型安全,避免手动转换的错误。
方案对比
- 方案1:无需生成额外类,灵活度高,适合Schema频繁变动的场景;但需要手动处理所有字段,代码量较大。
- 方案2:类型安全,代码更简洁,适合Schema固定的生产环境;但需要依赖Avro工具生成类,增加构建步骤。
内容的提问来源于stack exchange,提问作者Gabber
相关产品推荐
相关产品推荐

