尧图网络科技YAOTU DIGITAL 获取报价
获取报价
首页 / 资讯中心 / 文章详情

Apache Beam Java 读取 Parquet 文件实战:ParquetIO.read 源码级拆解与完整示例

发布时间:2026/9/28 2:35:31

资讯中心
01
ARTICLE

Apache Beam Java 读取 Parquet 文件实战:ParquetIO.read 源码级拆解与完整示例

Apache Beam Java 读取 Parquet 文件实战:ParquetIO.read 源码级拆解与完整示例
大数据批处理流处理数据工程【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam4/beam点击查看免费下载本篇技术指南聚焦 Apache Beam 官方代码解释文档中ReadParquetFile示例learning/prompts/code-explanation/java/10_io_parquet.md完整讲解如何用 Beam Java SDK 的ParquetIO连接器从 Parquet 文件读取数据覆盖反射式 Avro Schema 定义、PipelineOptions 命令行参数解析、ParquetIO.read与ParDo的组装并结合 ParquetIO.java 源码与 ParquetIOTest.java 测试用例深入其 Splittable DoFn 并行读取、Avro 数据模型、列投影、Beam Schema 推断等底层机制。读完本文你将能独立编写、运行并调试一个基于 Parquet 的 Beam 批处理读取管线并理解如何扩展到readFiles、parseGenericRecords和 Parquet 写入等进阶场景。一、程序全景一个最小可运行的 Parquet 读取管线示例代码ReadParquetFile是一个完整的 Beam 批处理程序其整体执行流程为通过PipelineOptionsFactory解析命令行参数构建ReadParquetFileOptions用ReflectData.get().getSchema(...)从 POJO 反射出 AvroSchemaPipeline.create(options)创建管线ParquetIO.read(schema).withAvroDataModel(GenericData.get()).from(filePattern)读取 Parquet 文件产出PCollectionGenericRecordParDo中DoFn逐条取出GenericRecord字段并输出日志p.run()触发执行。完整代码摘自原文档可直接编译运行package parquet; import java.io.Serializable; import java.util.Objects; import org.apache.avro.Schema; import org.apache.avro.generic.GenericData; import org.apache.avro.generic.GenericRecord; import org.apache.avro.reflect.ReflectData; import org.apache.beam.sdk.Pipeline; import org.apache.beam.sdk.io.parquet.ParquetIO; import org.apache.beam.sdk.options.Description; import org.apache.beam.sdk.options.PipelineOptions; import org.apache.beam.sdk.options.PipelineOptionsFactory; import org.apache.beam.sdk.options.Validation; import org.apache.beam.sdk.transforms.DoFn; import org.apache.beam.sdk.transforms.ParDo; import org.slf4j.Logger; import org.slf4j.LoggerFactory; public class ReadParquetFile { private static final Logger LOG LoggerFactory.getLogger(ReadParquetFile.class); public static class ExampleRecord implements Serializable { public static final String COLUMN_ID id; public static final String COLUMN_MONTH month; public static final String COLUMN_AMOUNT amount; private int id; private String month; private String amount; } public interface ReadParquetFileOptions extends PipelineOptions { Description(A glob file pattern to read Parquet files from) Validation.Required String getFilePattern(); void setFilePattern(String filePattern); } public static void main(String[] args) { ReadParquetFileOptions options PipelineOptionsFactory.fromArgs(args).withValidation().as(ReadParquetFileOptions.class); Schema exampleRecordSchema ReflectData.get().getSchema(ExampleRecord.class); Pipeline p Pipeline.create(options); p.apply( Read from Parquet file, ParquetIO.read(exampleRecordSchema) .withAvroDataModel(GenericData.get()) .from(options.getFilePattern())) .apply( Log records, ParDo.of( new DoFnGenericRecord, GenericRecord() { ProcessElement public void processElement(ProcessContext c) { GenericRecord record Objects.requireNonNull(c.element()); LOG.info( Id {}, Name {}, Amount {}, record.get(ExampleRecord.COLUMN_ID), record.get(ExampleRecord.COLUMN_MONTH), record.get(ExampleRecord.COLUMN_AMOUNT)); c.output(record); } })); p.run(); } }运行方式使用 DirectRunner 本地调试# 参数名来自 getFilePattern()/setFilePattern()即 --filePattern ./gradlew :sdks:java:io:parquet:run \ --args--runnerDirectRunner --filePattern/data/events/*.parquet需要说明原文档的英文解释中将该命令行参数写作--path但根据代码中getFilePattern()/setFilePattern()的命名PipelineOptionsFactory实际生成的命令行参数是--filePattern取值支持 glob 通配符。这一点以代码实现为准。二、用反射定义 Avro SchemaExampleRecord 与 ReflectDataExampleRecord是承载字段名的 POJO同时是 Avro Schema 的来源。它实现了Serializable因为 Beam 中的DoFn、管道选项等组件需要在执行器之间序列化传输字段名常量COLUMN_ID、COLUMN_MONTH、COLUMN_AMOUNT与 Parquet 文件中的列名一一对应public static class ExampleRecord implements Serializable { public static final String COLUMN_ID id; public static final String COLUMN_MONTH month; public static final String COLUMN_AMOUNT amount; private int id; private String month; private String amount; }Schema 的生成由 Avro 的ReflectData完成Schema exampleRecordSchema ReflectData.get().getSchema(ExampleRecord.class);ReflectData.get()会通过 Java 反射读取类的字段id、month、amount将其转换为等价的 Avro record Schemaid映射为 int、month/amount映射为 string。ParquetIO.read要求显式提供 Schema因为它需要与文件 footer 中存储的 Parquet 消息类型对齐才能正确反序列化列数据。这里有一个容易被忽略的细节被仓库测试用例明确锁定用反射生成的 Schema 读取时必须配合withAvroDataModel(GenericData.get())。在 ParquetIOTest.java 中testWriteAndReadUsingReflectDataSchemaWithoutDataModelThrowsException用ReflectDataSchema 读取但不指定数据模型直接抛出PipelineExecutionExceptiontestWriteAndReadUsingReflectDataSchemaWithDataModel同一场景加上withAvroDataModel(GenericData.get())后读写全部成功。原因在于ReflectData生成的 Schema 带有面向反射类的类型映射若不显式指定GenericData数据模型Parquet 的AvroParquetReader会尝试按反射模型实例化记录导致类型不匹配。这一点正是示例代码在第 57 行调用.withAvroDataModel(GenericData.get())的用意见下一节。三、用 PipelineOptions 定义命令行参数ReadParquetFileOptions展示了 Beam 标准的管道选项Pipeline Options模式一个继承PipelineOptions的接口通过 getter/setter 声明参数public interface ReadParquetFileOptions extends PipelineOptions { Description(A glob file pattern to read Parquet files from) Validation.Required String getFilePattern(); void setFilePattern(String filePattern); }Description生成帮助信息时显示的参数说明Validation.Required配合withValidation()在启动时强制校验若未提供--filePattern会在管线创建前报错避免把错误延迟到运行期getter/setter 命名直接决定命令行参数名getFilePattern()/setFilePattern()对应--filePatternglob。参数解析发生在main方法的开头ReadParquetFileOptions options PipelineOptionsFactory.fromArgs(args).withValidation().as(ReadParquetFileOptions.class);PipelineOptionsFactory.fromArgs(args)解析命令行withValidation()触发Validation.Required等注解校验as(ReadParquetFileOptions.class)将通用PipelineOptions转换为自定义子接口。随后Pipeline.create(options)将解析好的选项注入管线Runner 选择、文件路径等都从这份 options 中读取。四、ParquetIO.read连接器核心与 Splittable DoFn 并行机制读取的核心是ParquetIO.readp.apply( Read from Parquet file, ParquetIO.read(exampleRecordSchema) .withAvroDataModel(GenericData.get()) .from(options.getFilePattern()))在 ParquetIO.java 中read(Schema)返回一个Read对象它继承PTransformPBegin, PCollectionGenericRecord即从管道起点读入、输出PCollectionGenericRecordpublic static Read read(Schema schema) { return new AutoValue_ParquetIO_Read.Builder() .setSchema(schema) .setInferBeamSchema(false) .build(); }值得注意的配置项from(String filepattern)接受 glob 通配符如hdfs://namenode/data/*.parquet、gs://bucket/path/*.parquet也提供接受ValueProviderString的重载以支持运行时模板参数withAvroDataModel(GenericData model)设置AvroParquetReader的数据模型对应 Parquet 的AvroParquetReader.Builder#withDataModel(GenericData)withProjection(Schema projectionSchema, Schema encoderSchema)只读取投影列提升读取效率withBeamSchemas(boolean)开启 Beam Schema 推断使结果可用于 Beam SQL 等 Schema 化操作withConfiguration(MapString, String | Configuration)透传 HadoopConfiguration如设置 Parquet 过滤谓词、读取参数。Read.expand的实现展示了它的底层链路ParquetIO.javaPCollectionReadableFile inputFiles input .apply(Create filepattern, Create.ofProvider(getFilepattern(), StringUtf8Coder.of())) .apply(FileIO.matchAll()) .apply(FileIO.readMatches()); ... return inputFiles.apply(readFiles);也就是说ParquetIO.read本质上是FileIO.matchAll/readMatches与readFiles的组合。正因为底层是FileIO它天然支持任何 Beam 文件系统HDFS、GCS、S3、本地文件系统等。真正的逐块读取由内部类SplitReadFn完成——这是一个Splittable DoFnSDFParquetIO.javaGetInitialRestriction打开文件 footer把行组Row Group数量作为初始偏移范围OffsetRange(0, rowGroups.size())SplitRestriction按SPLIT_LIMIT 6400000064 MB把行组聚合成可并行处理的块即先把文件切成 64MB 大小的块运行时可进一步动态切分以提高读取并行度ProcessElement使用ParquetFileReader逐行组读取tracker.tryClaim(currentBlock)驱动动态分块dynamic splitting行组内逐条recordReader.read()并调用parseFn输出BlockTracker extends OffsetRangeTracker负责进度与分裂汇报进度按平均记录大小估算ParquetIOTest.java 的testBlockTracker验证了其工作完成量从 0 到 1 的推进BeamParquetInputFile把 Beam 的SeekableByteChannel适配为 Parquet 的InputFile使得 Parquet 读取逻辑完全复用 parquet-mr 的实现。因此ParquetIO.read不是简单的逐文件顺序读而是基于 Splittable DoFn 的文件→64MB 块→行组→记录四级并行分解这是它在大数据集上获得良好伸缩性的关键。testSplitBlockWithLimitParquetIOTest.java则精确验证了按字节上限切分行组的算法。五、withAvroDataModelAvro 数据模型与兼容性开关withAvroDataModel(GenericData.get())把读取使用的数据模型固定为GenericData。在SplitReadFn.getConfWithModelClass()ParquetIO.java中可以看到它与 Hadoop 配置的联动if (model ! null (model.getClass() GenericData.class || model.getClass() SpecificData.class)) { conf.setBoolean(AvroReadSupport.AVRO_COMPATIBILITY, true); } else { conf.setBoolean(AvroReadSupport.AVRO_COMPATIBILITY, false); }指定GenericData通用记录模型或SpecificData特定类模型时开启AVRO_COMPATIBILITYavro.compatibility走标准的 Avro 兼容读取路径数据模型为null时即不调用withAvroDataModel该标志关闭反射 Schema 场景下会与写入端模型不一致而抛异常正是 ParquetIOTest.java 中testWriteAndReadUsingReflectDataSchemaWithoutDataModelThrowsException所验证的行为。实践建议只要 Schema 来自ReflectData即 POJO 反射生成读取端就应像示例代码一样显式调用withAvroDataModel(GenericData.get())若 Schema 是手工构造或从已有 Avro Schema 文件加载则不是必须的。六、ParDo 处理与日志输出读取后的PCollectionGenericRecord交给ParDo逐元素处理.apply( Log records, ParDo.of( new DoFnGenericRecord, GenericRecord() { ProcessElement public void processElement(ProcessContext c) { GenericRecord record Objects.requireNonNull(c.element()); LOG.info( Id {}, Name {}, Amount {}, record.get(ExampleRecord.COLUMN_ID), record.get(ExampleRecord.COLUMN_MONTH), record.get(ExampleRecord.COLUMN_AMOUNT)); c.output(record); } }));要点DoFnGenericRecord, GenericRecord的输入输出类型都是GenericRecordProcessElement中的ProcessContext提供element()当前记录与output()下游发射两个能力Objects.requireNonNull(c.element())对空元素做了防御性校验——Parquet 读取端对被过滤掉/损坏被跳过的记录本就不会输出此处的空值防御可进一步避免静默空指针按record.get(字段名常量)取列值。需注意示例日志文本写作Id {}, Name {}, Amount {}而实际字段名是id、month、amount日志文案与列名不完全一致这是示例代码本身的小瑕疵可自行对齐为Id {}, Month {}, Amount {}c.output(record)将记录原样透传使得后续还可以继续挂接其他变换如ParquetIO的写入端本示例到日志为止。Pipeline.run()将整个 DAG 提交执行使用 DirectRunner 时同步执行使用 Dataflow 等分布式 Runner 时提交远程作业。七、进阶扩展readFiles、未知 Schema 解析与 Beam SchemaParquetIO不止提供read围绕文件来源与 Schema 是否已知提供了完整的能力矩阵见 ParquetIO.java 与 ParquetIOTest.java 的对应测试场景变换说明Schema 已知文件来自 globParquetIO.read(schema).from(pattern)本示例使用基于 FileIOSchema 已知文件来自PCollectionFileIO.ReadableFileParquetIO.readFiles(schema)更灵活可与FileIO.matchAll/readMatches自由组合Schema 未知按 glob 读取ParquetIO.parseGenericRecords(parseFn).from(pattern)需要提供SerializableFunctionGenericRecord, T解析函数将记录转为自定义类型Schema 未知文件来自PCollectionReadableFileParquetIO.parseFilesGenericRecords(parseFn)上述场景的文件集合版本需要 SQL / Schema 化操作read(...).withBeamSchemas(true)自动推断 Beam Schema 并安装对应 CoderSchemaCoder典型用法来自源码 Javadoc 与测试// 1) 从文件集合读取 PCollectionFileIO.ReadableFile files pipeline .apply(FileIO.match().filepattern(options.getInputFilepattern())) .apply(FileIO.readMatches()); PCollectionGenericRecord output files.apply(ParquetIO.readFiles(SCHEMA)); // 2) 未知 Schema转换为自定义类型甚至转为 Beam Row 使用 SQL PCollectionString jsonRecords p.apply( ParquetIO.parseGenericRecords(record - record.toString()).from(*.parquet));关于编码器Coder默认情况下ReadFiles使用AvroCoder.of(schema)开启withBeamSchemas(true)后改用AvroUtils.schemaCoder(schema)即SchemaCodergetCollectionCoder()的实现见 ParquetIO.java。测试testWriteAndReadWithBeamSchemaParquetIOTest.java验证了该模式下读写数据完全一致。另外parseFilesGenericRecords内部通过TypeDescriptors.outputOf(getParseFn())判断输出类型若解析函数直接输出GenericRecord会抛出Parse cant be used for reading as GenericRecord.ParquetIOTest.java 的testReadFilesUnknownSchemaFilesForGenericRecordThrowException验证了这一点——因为未知 Schema 场景本就无法确定记录结构必须转换为明确类型Coder 优先取显式withCoder(...)否则从CoderRegistry按解析函数输出类型推断。列投影读取Projection当只需要部分列时可用投影 Schema 减少处理与编码开销PCollectionGenericRecord records p.apply( ParquetIO.read(SCHEMA) .from(/foo/bar/*.parquet) .withProjection(projectionSchema, encoderSchema));projectionSchema声明要读取的列encoderSchema声明输出编码时使用的 Schema未请求的列置为 nullable。源码 Javadoc 指出投影读取会自动启用 Splittable 读取虽然能节省读取/编码时间但收益并不与请求的数据量成正比——因为 Parquet 按列交错存储读取器仍需按编码 Schema 遍历数据集节省的只是读取不需要列的时间。测试testWriteAndReadWithProjectionParquetIOTest.java验证了投影读取结果与预期一致。Hadoop 配置透传与谓词下推withConfiguration(...)允许把 HadoopConfiguration传给 Parquet 读取器例如在读取端设置ParquetInputFormat.setFilterPredicate做谓词下推ParquetIOTest.java 的testWriteAndReadWithConfiguration演示了按id 0过滤后只读出 1 条记录。populateDisplayData会把所有以parquet开头的配置项登记到 Beam 的 Display Data 中便于 UI/监控查看ParquetIOTest.java 的testReadDisplayData。八、配套写入端与依赖配置ParquetIO.Sink写出 Parquet 文件虽然示例只读取但ParquetIO是完整的读写连接器。sink(Schema)返回FileIO.SinkGenericRecord配合FileIO.write使用其默认参数与可调项在 ParquetIO.java 的工厂方法中一目了然public static Sink sink(Schema schema) { return new AutoValue_ParquetIO_Sink.Builder() .setJsonSchema(schema.toString()) .setCompressionCodec(CompressionCodecName.SNAPPY) .setRowGroupSize(ParquetWriter.DEFAULT_BLOCK_SIZE) .setPageSize(ParquetWriter.DEFAULT_PAGE_SIZE) .setEnableDictionary(ParquetWriter.DEFAULT_IS_DICTIONARY_ENABLED) .setEnableBloomFilter(ParquetProperties.DEFAULT_BLOOM_FILTER_ENABLED) .build(); }常用配置项汇总配置默认值作用withCompressionCodec(CompressionCodecName)SNAPPY压缩算法可选UNCOMPRESSED、GZIP、LZO、ZSTD等withRowGroupSize(int)ParquetWriter.DEFAULT_BLOCK_SIZE行组大小必须为正数withPageSize(int)1 MBDEFAULT_PAGE_SIZE列页大小withDictionaryEncoding(boolean)默认开启字典编码开关withBloomFilterEnabled(boolean)默认关闭Bloom 过滤器开关withMinRowCountForPageSizeCheck(int)100大行场景可调小如 1页大小检查前最少缓冲行数避免大行导致内存暴涨withMaxRowCountForPageSizeCheck(int)估计机制可延迟至 10000 行强制页大小检查的最大缓冲行数上限配合上行可约束行宽差异大的表withAvroDataModel(GenericData)null写入端数据模型写入示例p.apply(...) // PCollectionGenericRecord .apply(FileIO.GenericRecordwrite() .via(ParquetIO.sink(SCHEMA) .withCompressionCodec(CompressionCodecName.SNAPPY)) .to(destination/path) .withSuffix(.parquet));Sink.open(WritableByteChannel)内部把 Beam 通道包装成BeamParquetOutputFile实现 Parquet 的OutputFile接口随后构建AvroParquetWriterflush()通过writer.close()完成真正落盘ParquetIO.java。模块依赖与版本ParquetIO 模块位于 sdks/java/io/parquet/构建文件 build.gradle 显示核心依赖parquet-avro、parquet-column、parquet-common、parquet-hadoop版本统一为1.15.2依赖 Beam:sdks:java:coreshadow 配置与:sdks:java:extensions:avro、:sdks:java:io:hadoop-commonHadoop 以provided方式引入hadoop-client、hadoop-common并内置了 2.10.2 / 3.2.4 / 3.3.6 / 3.4.1 多版本兼容测试任务hadoopVersionXxxTest验证不同 Hadoop 版本下的读取行为。九、测试验证与调试建议仓库为ParquetIO提供了完整的测试覆盖ParquetIOTest.java共 644 行可作为自行调试时的对照基准testWriteAndRead写入 1000 条GenericRecord再读回用PAssert.that(readBack).containsInAnyOrder(records)断言一致性testWriteWithRowGroupSizeAndRead验证自定义行组大小不影响读写一致性testWriteAndReadWithProjection验证投影 Schema 读回记录与预期一致testWriteAndReadWithBeamSchema验证withBeamSchemas(true)下的完整读写testWriteAndReadAsJsonForUnknownSchema/testReadFilesAsRowForUnknownSchemaFiles验证未知 Schema 场景下的parseGenericRecords/parseFilesGenericRecords与 BeamRow转换testWriteAndReadWithConfiguration验证 Hadoop 配置透传与谓词过滤下推集成测试见 ParquetIOIT.java。调试ReadParquetFile时的常见问题与对策--filePattern缺失报错Validation.Required会在启动时拦截确认命令行参数名与 getter 命名一致反射 Schema 读取抛PipelineExecutionException补齐.withAvroDataModel(GenericData.get())Coder 推断失败parseGenericRecords输出自定义类型时显式.withCoder(...)并行度不达预期检查行组分布与 64MB 块切分SPLIT_LIMIT小文件多时并行度受文件数量限制日志看不到输出确认使用 SLF4J 配置本地可加slf4j-simple依赖并注意 DirectRunner 下p.run()同步执行完成后再查看日志。总结ReadParquetFile示例虽短却串起了 Apache Beam ParquetIO 连接器的核心链路反射式 Avro Schema →PipelineOptions参数化 →ParquetIO.read的 Splittable DoFn 并行读取 →ParDo处理 →Pipeline.run执行。在此基础上readFiles、parseGenericRecords、withBeamSchemas、withProjection、withConfiguration与ParquetIO.sink构成了完整的读写能力矩阵。理解源码层文件→64MB 块→行组→记录的分解机制与 Avro 数据模型兼容性细节将帮助你在实际项目中正确配置参数、规避反射 Schema 陷阱并最大化 Parquet 大数据读取的性能与伸缩性。赞分享大数据批处理流处理数据工程【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam4/beam点击查看免费下载相关推荐Apache Beam Java Kata 实战用 TextIO.read() 从文本文件读取 PCollectionApache Beam Java Kata 实战用 TextIO.read 从文本文件读取 PCollection 本篇技术指南以 Apache Beam 官大数据批处理流处理数据工程Apache Arrow C 实战用 parquet/arrow 接口读写 Parquet 文件附完整示例工程解析Apache Arrow C 实战用 parquet/arrow 接口读写 Parquet 文件附完整示例工程解析 本指南以 Apache Arrow数据工程大数据序列化数据分析Apache Beam Java Cookbook六大常用数据分析模式的示例实现与源码解读Apache Beam Java Cookbook六大常用数据分析模式的示例实现与源码解读 本文围绕 Apache Beam 仓库中 examples/jav大数据批处理流处理数据工程上一篇Awaken全平台EPUB阅读器终极指南打造你的个人数字图书馆下一篇告别选择困难2025最值得尝试的5个Arduino/Raspberry Pi开源项目创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
02
RELATED NEWS

相关资讯

更多网站建设与数字化升级内容

03
WHY YAOTU

想打造同款高转化官网?

懂行业、懂生意,从建站到增长一站式陪跑

◈

场景化定制

不做模板站,围绕你的业务场景量身设计,小众不撞款。

◐

营销型架构

以转化目标组织内容与路径,让官网真正带来询盘。

▲

全周期服务

设计、开发、运营、运维一体,上线只是开始。

免费获取你的建站方案

留下需求,专属顾问 24 小时内为你输出方案建议。