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

如何用Apache Beam Java实现两个CSV文件的ID匹配对比?

Apache Beam Java 实现CSV文件ID匹配方案

原代码核心问题

  1. 违反Beam编程模型:在DoFn的processElement方法内对PCollection调用apply是完全错误的——DoFn仅用于处理单个元素,无法直接操作分布式数据集PCollection。
  2. 匹配逻辑不符合需求:当前代码是整行相等匹配,而非需求中的ID字段匹配,缺少CSV行解析提取ID的步骤。

正确实现方案

根据需求,我们需要将file1的每行ID与file2的所有ID匹配,保留匹配成功的行。推荐两种实现方式,可根据数据集大小选择:

方案一:侧输入(Side Input)适合小数据集匹配

当其中一个文件(如file2)较小时,将其ID集合作为侧输入传递给处理file1的DoFn,效率更高。

完整代码

import org.apache.beam.sdk.Pipeline;
import org.apache.beam.sdk.io.TextIO;
import org.apache.beam.sdk.transforms.DoFn;
import org.apache.beam.sdk.transforms.ParDo;
import org.apache.beam.sdk.transforms.View;
import org.apache.beam.sdk.values.PCollection;
import org.apache.beam.sdk.values.PCollectionView;
import org.apache.commons.csv.CSVFormat;
import org.apache.commons.csv.CSVParser;
import org.apache.commons.csv.CSVRecord;

import java.io.StringReader;
import java.util.Set;

public class CsvIdMatcher {
    public static void main(String[] args) {
        Pipeline pipeline = Pipeline.create();

        // 读取CSV文件,假设第一列是ID
        PCollection<String> file1 = pipeline.apply(TextIO.read().from("gs://my-bucket/file1.csv"));
        PCollection<String> file2 = pipeline.apply(TextIO.read().from("gs://my-bucket/file2.csv"));

        // 处理file2:提取所有ID并转换成侧输入(Set集合)
        PCollectionView<Set<String>> file2Ids = file2
                .apply(ParDo.of(new ExtractIdFn()))
                .apply(View.asSet());

        // 处理file1:检查每行ID是否在file2的ID集合中,匹配则输出整行
        PCollection<String> matchedRows = file1.apply(ParDo.of(new MatchIdFn(file2Ids)));

        // 写入结果文件
        matchedRows.apply(TextIO.write()
                .to("gs://my-bucket/comparison_results.csv")
                .withSuffix(".csv")
                .withNumShards(1)); // 可选:合并成单个文件

        pipeline.run().waitUntilFinish();
    }

    // 提取CSV行中的ID(假设ID是第一列)
    public static class ExtractIdFn extends DoFn<String, String> {
        @ProcessElement
        public void processElement(@Element String line, OutputReceiver<String> out) {
            try (CSVParser parser = CSVParser.parse(line, CSVFormat.DEFAULT.withFirstRecordAsHeader(false))) {
                for (CSVRecord record : parser) {
                    String id = record.get(0); // 替换为实际ID列的索引或字段名
                    out.output(id.trim());
                }
            } catch (Exception e) {
                // 处理解析错误,可跳过错误行或记录日志
                System.err.println("Failed to parse line: " + line + " Error: " + e.getMessage());
            }
        }
    }

    // 匹配ID并输出整行
    public static class MatchIdFn extends DoFn<String, String> {
        private final PCollectionView<Set<String>> file2Ids;

        public MatchIdFn(PCollectionView<Set<String>> file2Ids) {
            this.file2Ids = file2Ids;
        }

        @ProcessElement
        public void processElement(@Element String line, OutputReceiver<String> out, @SideInput PCollectionView<Set<String>> file2Ids) {
            try (CSVParser parser = CSVParser.parse(line, CSVFormat.DEFAULT.withFirstRecordAsHeader(false))) {
                for (CSVRecord record : parser) {
                    String id = record.get(0); // 替换为实际ID列的索引或字段名
                    if (file2Ids.get().contains(id.trim())) {
                        out.output(line); // 输出file1的整行,若需要同时输出file2的行可改用Join方案
                    }
                }
            } catch (Exception e) {
                System.err.println("Failed to parse line: " + line + " Error: " + e.getMessage());
            }
        }
    }
}

关键说明

  • 使用View.asSet()将file2的ID集合转换为侧输入,允许DoFn高效访问整个数据集的ID集合。
  • 使用@SideInput注解在DoFn中获取侧输入。
  • 引入Apache Commons CSV库解析CSV行,需在pom.xml中添加依赖:
    <dependency>
        <groupId>org.apache.commons</groupId>
        <artifactId>commons-csv</artifactId>
        <version>1.10.0</version>
    </dependency>
    

方案二:CoGroupByKey 实现分布式Join(适合大数据集)

当两个文件都较大时,使用Beam的CoGroupByKey进行分布式Join,按ID分组后处理匹配结果。

核心代码片段

import org.apache.beam.sdk.transforms.CoGroupByKey;
import org.apache.beam.sdk.values.CoGbkResult;
import org.apache.beam.sdk.values.KV;
import org.apache.beam.sdk.values.KeyedPCollectionTuple;

// ... 读取文件部分同方案一

// 将file1转换为KV<ID, 行>
PCollection<KV<String, String>> file1Kv = file1.apply(ParDo.of(new ToKvFn()));
// 将file2转换为KV<ID, 行>
PCollection<KV<String, String>> file2Kv = file2.apply(ParDo.of(new ToKvFn()));

// 按ID分组
PCollection<KV<String, CoGbkResult>> grouped = KeyedPCollectionTuple
        .of("file1", file1Kv)
        .and("file2", file2Kv)
        .apply(CoGroupByKey.create());

// 处理分组结果,输出匹配的行
PCollection<String> matchedRows = grouped.apply(ParDo.of(new ProcessJoinedFn()));

// ... 写入结果部分同方案一

// 辅助DoFn:将CSV行转换为KV<ID, 行>
public static class ToKvFn extends DoFn<String, KV<String, String>> {
    @ProcessElement
    public void processElement(@Element String line, OutputReceiver<KV<String, String>> out) {
        try (CSVParser parser = CSVParser.parse(line, CSVFormat.DEFAULT.withFirstRecordAsHeader(false))) {
            for (CSVRecord record : parser) {
                String id = record.get(0).trim();
                out.output(KV.of(id, line));
            }
        } catch (Exception e) {
            System.err.println("Failed to parse line: " + line + " Error: " + e.getMessage());
        }
    }
}

// 处理分组后的匹配结果
public static class ProcessJoinedFn extends DoFn<KV<String, CoGbkResult>, String> {
    @ProcessElement
    public void processElement(@Element KV<String, CoGbkResult> element, OutputReceiver<String> out) {
        String id = element.getKey();
        Iterable<String> file1Rows = element.getValue().getAll("file1");
        Iterable<String> file2Rows = element.getValue().getAll("file2");
        
        // 输出所有匹配的组合,可根据需求调整格式
        for (String f1Row : file1Rows) {
            for (String f2Row : file2Rows) {
                out.output(String.format("File1: %s\nFile2: %s\n", f1Row, f2Row));
            }
        }
    }
}

关键说明

  • 将两个文件都转换为KV<ID, 行>格式,通过CoGroupByKey按ID分组。
  • 分组后,每个Key对应的CoGbkResult包含所有匹配的行,可按需输出组合结果。

内容的提问来源于stack exchange,提问作者sidd

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.14 06:16:01