使用Avro接口读取Parquet文件中深度嵌套记录失败求助
读取嵌套数组+多层记录结构的Parquet数据(Parquet-MR 1.8.1)
刚好我之前用Parquet-MR 1.8.1处理过几乎一模一样的嵌套结构数据,给你捋捋具体的读取步骤和代码示例,绝对实用!
先把你给出的Schema整理得更清晰些,方便对照:
{ "type": "record", "name": "record", "fields": [ { "name": "elements", "type": { "type": "array", "items": { "type": "record", "name": "elementWrapper", "fields": [ { "name": "array_element", "type": { "type": "record", "name": "element", "namespace": "test", "fields": [ { "name": "someField", "type": "int" } ] } } ] } } } ] }
第一步:生成对应Java实体类(基于Avro Schema)
Parquet-MR和Avro的适配性最好,所以第一步我们用Avro工具把Schema转成Java类,这样后续读取时能直接映射到对象,不用手动解析字段。
如果你用Maven的话,直接在pom.xml里加这个插件配置:
<plugin> <groupId>org.apache.avro</groupId> <artifactId>avro-maven-plugin</artifactId> <version>1.7.7</version> <!-- 这个版本和Parquet-MR 1.8.1完美兼容,别乱换 --> <executions> <execution> <phase>generate-sources</phase> <goals> <goal>schema</goal> </goals> <configuration> <sourceDirectory>${project.basedir}/src/main/resources/</sourceDirectory> <outputDirectory>${project.basedir}/src/main/java/</outputDirectory> </configuration> </execution> </executions> </plugin>
把你的Schema保存成record.avsc文件放在src/main/resources目录下,然后执行mvn generate-sources命令,就能自动生成Record.java、ElementWrapper.java和test.Element.java这几个实体类,字段和嵌套结构完全对应。
第二步:用ParquetAvroReader读取并解析数据
接下来写读取代码,核心思路就是逐层遍历:先拿外层的elements数组,再遍历每个数组元素(也就是elementWrapper),最后取出里面的element对象读取具体字段。
示例代码如下:
import org.apache.parquet.avro.AvroParquetReader; import org.apache.parquet.hadoop.ParquetReader; import org.apache.parquet.hadoop.util.HadoopInputFile; import org.apache.hadoop.fs.Path; import test.Element; import java.io.IOException; public class ParquetNestedDataReader { public static void main(String[] args) throws IOException { // 替换成你的Parquet文件路径 Path parquetFilePath = new Path("/your/parquet/file/path/xxx.parquet"); // 初始化ParquetReader,指定我们生成的Record类作为泛型 ParquetReader<Record> reader = AvroParquetReader.<Record>builder(HadoopInputFile.fromPath(parquetFilePath)) .build(); Record currentRecord; // 逐行读取Parquet文件 while ((currentRecord = reader.read()) != null) { // 先判断数组是否为空,避免空指针 if (currentRecord.getElements() != null) { // 遍历elements数组里的每个elementWrapper for (ElementWrapper wrapper : currentRecord.getElements()) { // 取出包装类里的element对象 Element nestedElement = wrapper.getArrayElement(); // 读取具体字段值,这里是someField System.out.println("读取到someField值:" + nestedElement.getSomeField()); } } } // 记得关闭资源 reader.close(); } }
几个关键注意事项(针对Parquet-MR 1.8.1)
- 版本兼容性:一定要用Avro 1.7.7,和Parquet-MR 1.8.1的版本完全匹配,不然很容易出现序列化/反序列化的异常,踩过坑的过来人提醒你!
- 空值处理:实际业务里数组或者嵌套字段可能为空,一定要加非空判断,不然运行时会抛出空指针异常。
- 性能优化:如果处理的是超大Parquet文件,可以考虑调整ParquetReader的配置,比如设置
withConf()传入Hadoop配置,调整缓冲区大小,或者用批量读取的方式提升效率。
内容的提问来源于stack exchange,提问作者Iain
相关产品推荐
相关产品推荐

