如何使用Alpakka将Parquet存储的JSON记录映射为Scala样例类
问题分析
你的报错本质是AvroParquetSource要求输入的ParquetReader产出的类型必须是org.apache.avro.generic.GenericRecord的子类,你自定义的Person样例类和Avro的GenericRecord没有继承关系,所以直接替换类型会触发泛型边界校验失败。
解决方案
有两种常用的实现路径,可根据你的业务场景选择:
方案一:在读取流中增加GenericRecord到Person的映射步骤
这是改造成本最低的方案,不需要修改原有Parquet读取逻辑,只需要在Source后追加转换步骤即可:
import org.apache.avro.generic.GenericRecord // 原有读取逻辑保持不变 val reader: ParquetReader[GenericRecord] = AvroParquetReader.builder[GenericRecord](filePath).withConf(conf).build() val source: Source[GenericRecord, NotUsed] = AvroParquetSource(reader) // 新增类型转换逻辑 val personSource: Source[Person, NotUsed] = source.map { record => Person( name = record.get("name").toString, age = record.get("age").asInstanceOf[Double] ) } // 直接使用转换后的流调用ask方法 personSource.ask[WorkerAck](28)(workerActor)
注意:如果你的Parquet字段允许为空,需要提前添加判空逻辑避免运行时空指针
方案二:使用Avro规范编译生成的样例类(适合字段多、结构复杂的场景)
如果你需要更高的类型安全性,可以通过Avro Schema编译生成原生实现GenericRecord接口的Person类,直接满足AvroParquetSource的泛型要求:
- 先定义Person的Avro Schema文件
person.avsc:
{ "type": "record", "name": "Person", "namespace": "com.common", "fields": [ {"name": "name", "type": "string"}, {"name": "age", "type": "double"} ] }
- 通过avro-maven-plugin(Maven项目)或sbt-avro插件(SBT项目)在编译阶段自动生成Person类
- 直接用特定类型读取Parquet文件:
import org.apache.parquet.avro.AvroParquetReader import org.apache.avro.generic.GenericData import com.common.Person val reader: ParquetReader[Person] = AvroParquetReader.builder[Person](filePath) .withConf(conf) .withModel(GenericData.get()) .build() val source: Source[Person, NotUsed] = AvroParquetSource(reader) source.ask[WorkerAck](28)(workerActor)
该方案不需要手动写字段映射逻辑,编译期即可校验字段类型匹配性,避免运行时转换错误。
内容的提问来源于stack exchange,提问作者igx
相关产品推荐
相关产品推荐

