You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.09.26 18:15:11