如何在Apache Beam中对PCollection追加行或执行Union合并操作
实现方案
你当前的RowAddition实现逻辑存在问题:它会对每一条原始行都输出一条待追加的新记录,最终会生成6条重复的新增行,不符合预期。无需复杂逻辑,按以下步骤调整即可:
1. 构造待追加行的独立PCollection
直接在管道中创建仅包含你要追加的那一行记录的PCollection,和原始读取到的CSV行PCollection结构完全一致。
2. 用Flatten合并两个同结构PCollection
Apache Beam中合并多个结构完全相同的PCollection,使用Flatten转换即可,无需使用关联操作。
完整修改代码
public static void main(String[] args) throws IOException { final File schemaFile = new File("src/main/resources/addRow/schema_transform.avsc"); File csvFile = new File("src/main/resources/addRow/CustomerRequest.csv"); Schema avroSchema = new Schema.Parser().parse(schemaFile); Pipeline pipeline = Pipeline.create(); org.apache.beam.sdk.schemas.Schema beamSchema = AvroUtils.toBeamSchema(avroSchema); final PCollectionTuple tuples = pipeline .apply("匹配CSV文件", FileIO.match().filepattern(csvFile.getAbsolutePath())) .apply("读取匹配到的文件", FileIO.readMatches()) .apply("解析校验为行对象", ParDo.of(new FileReader(beamSchema)) .withOutputTags(FileReader.validTag(), TupleTagList.of(invalidTag()))); // 读取到的原始有效CSV行 final PCollection<Row> originalRows = tuples.get(FileReader.validTag()) .setCoder(RowCoder.of(beamSchema)); // 构造待追加的新行 Row newAppendRow = Row.withSchema(beamSchema) .addValues("01", "30/07/2021", 999) .build(); // 生成仅包含待追加行的PCollection PCollection<Row> appendRowPcoll = pipeline.apply(Create.of(newAppendRow)) .setCoder(RowCoder.of(beamSchema)); // 合并原始行和待追加行 PCollection<Row> allRows = PCollectionList.of(originalRows) .and(appendRowPcoll) .apply(Flatten.pCollections()); // 转换为字符串输出为CSV PCollection<String> pOutput = allRows.apply(ParDo.of(new RowToString())); pOutput.apply(TextIO.write() .to("src/main/resources/addRow/rowOutput") .withNumShards(1) .withSuffix(".csv")); pipeline.run().waitUntilFinish(); System.out.println("执行结束"); }
原来的RowAddition类无需保留,可以直接删除。
注意事项
- 待追加行的日期格式要和原始数据保持一致,原始数据月日均为两位,因此填写
30/07/2021而非30/7/2021 - 两个需要合并的PCollection必须使用完全一致的Schema和Coder,否则会触发运行时报错
- 如果需要输出内容按日期排序,可以在合并后增加排序转换逻辑再输出
内容的提问来源于stack exchange,提问作者Ron Santis
相关产品推荐
相关产品推荐

