如何在Java Dataflow中读取多Schema Avro文件并保留GCS路径?
问题背景
我们原本使用Python实现了一个Dataflow任务,逻辑是监听Pub/Sub订阅获取GCS上Avro文件的路径(格式为gs://bucket/file-timestamp.avro),这些Avro文件没有统一的Schema。通过Python的avroio.ReadAllFromAvro(with_filename=True)可以将文件记录解析为字典,最终得到(GCS文件路径, 字典格式记录)的元组,Python代码如下:
# Python code pipeline | "Read from Pub/Sub" >> beam.io.ReadFromPubSub(subscription=input_subscription) | "Read Avro Files" >> avroio.ReadAllFromAvro(with_filename=True)
现在需要将该逻辑转为Java版本,已实现Pub/Sub读取,但无法在不指定固定Schema的情况下通过AvroIO得到PCollection<GenericRecord>,同时保留原始GCS文件路径。尝试结合FileIO和AvroIO的代码如下:
pipeline .apply("Read from Pub/Sub", PubsubIO.fromSubscription(options.getInputSubscription())) .apply("Read Avro Files", FileIO.matchAll()) .apply(FileIO.readMatches()) .apply(AvroIO.parseFilesGenericRecords(record -> record))
但触发报错:
java.lang.IllegalArgumentException: Unable to infer coder for output of parseFn. Specify it explicitly using withCoder()
由于Avro文件Schema不固定,无法指定特定Coder,需要可行的解决方案。
解决方案
要实现需求,需结合FileIO获取文件元数据(包含GCS路径),使用AvroIO解析为GenericRecord,同时通过通用Coder适配动态Schema,最终绑定文件路径与记录。以下提供两种可行实现方式:
方式一:手动在ParDo中解析并绑定路径
通过ParDo处理FileIO.ReadableFile,提取文件路径后读取Avro内容,使用AvroCoder.of(GenericRecord.class)作为通用Coder(自动从每个文件读取对应Schema):
import org.apache.beam.sdk.io.AvroIO; import org.apache.beam.sdk.io.FileIO; import org.apache.beam.sdk.io.gcp.pubsub.PubsubIO; import org.apache.beam.sdk.transforms.DoFn; import org.apache.beam.sdk.transforms.ParDo; import org.apache.avro.generic.GenericRecord; import org.apache.beam.sdk.coders.AvroCoder; import org.apache.beam.sdk.values.KV; // ... pipeline .apply("Read from Pub/Sub", PubsubIO.fromSubscription(options.getInputSubscription())) .apply("Match Avro Files", FileIO.matchAll()) .apply("Read File Matches", FileIO.readMatches()) .apply("Parse Avro with File Path", ParDo.of(new DoFn<FileIO.ReadableFile, KV<String, GenericRecord>>() { @ProcessElement public void processElement(@Element FileIO.ReadableFile file, OutputReceiver<KV<String, GenericRecord>> out) throws Exception { // 提取GCS文件路径 String filePath = file.getMetadata().resourceId().toString(); // 读取并解析Avro文件,使用通用Coder适配动态Schema try (AvroIO.Read.All<GenericRecord> reader = AvroIO.readGenericRecords(filePath) .withCoder(AvroCoder.of(GenericRecord.class))) { reader.read(file.open()) .forEach(record -> out.output(KV.of(filePath, record))); } } })) // 后续处理逻辑 // .apply(...)
方式二:使用AvroIO的简化API
利用AvroIO.parseFilesGenericRecords()的重载方法,直接指定输出类型为KV<String, GenericRecord>并绑定通用Coder,无需手动处理文件读取:
import org.apache.beam.sdk.io.AvroIO; import org.apache.beam.sdk.io.FileIO; import org.apache.beam.sdk.io.gcp.pubsub.PubsubIO; import org.apache.beam.sdk.values.KV; import org.apache.beam.sdk.values.TypeDescriptors; import org.apache.avro.generic.GenericRecord; import org.apache.beam.sdk.coders.AvroCoder; // ... pipeline .apply("Read from Pub/Sub", PubsubIO.fromSubscription(options.getInputSubscription())) .apply("Match Avro Files", FileIO.matchAll()) .apply("Parse Avro with Path", AvroIO.parseFilesGenericRecords() .withCoder(AvroCoder.of(GenericRecord.class)) .withOutputType(TypeDescriptors.kvs(TypeDescriptors.strings(), TypeDescriptors.of(GenericRecord.class)))) // 后续处理逻辑 // .apply(...)
关键说明
AvroCoder.of(GenericRecord.class)是核心:它不需要提前指定固定Schema,会自动从每个Avro文件中读取对应的Schema信息,适配动态Schema场景。- 两种方式最终均输出
KV<String, GenericRecord>,与Python版本的(gcs_filepath, record_as_dict)逻辑完全对应,可直接用于后续数据处理。
内容的提问来源于stack exchange,提问作者oikonomiyaki

