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

如何在Dataflow中将单条源记录生成两条含不同值的重复记录?

在Dataflow中实现单条记录拆分多条的方案

Python SDK 实现方式

方法1:使用FlatMap(推荐,代码更简洁)

FlatMap支持将单个输入元素转换为多个输出元素,直接定义映射函数即可实现拆分逻辑:

import apache_beam as beam

def split_record(record):
    # 假设record为字典格式,包含Column1字段
    if record['Column1'] == 'A':
        # 返回两条新记录,分别设置Column2的值
        return [
            {**record, 'Column2': 'A1'},
            {**record, 'Column2': 'A2'}
        ]
    # 非目标值时返回原记录(可根据实际需求调整逻辑)
    return [record]

with beam.Pipeline() as p:
    (p
     | '读取源数据' >> beam.io.ReadFromText('input.txt')  # 替换为实际数据源
     | '解析为字典' >> beam.Map(lambda line: dict(zip(['Column1'], line.split(','))))  # 按需修改解析逻辑
     | '拆分记录' >> beam.FlatMap(split_record)
     | '输出结果' >> beam.io.WriteToText('output.txt')
    )

方法2:使用自定义ParDo

通过继承DoFn,在process方法中多次yield输出多条记录:

import apache_beam as beam

class SplitRecordDoFn(beam.DoFn):
    def process(self, record):
        if record['Column1'] == 'A':
            yield {**record, 'Column2': 'A1'}
            yield {**record, 'Column2': 'A2'}
        else:
            yield record

with beam.Pipeline() as p:
    (p
     | '读取源数据' >> beam.io.ReadFromText('input.txt')
     | '解析为字典' >> beam.Map(lambda line: dict(zip(['Column1'], line.split(','))))
     | '拆分记录' >> beam.ParDo(SplitRecordDoFn())
     | '输出结果' >> beam.io.WriteToText('output.txt')
    )

Java SDK 实现方式

方法1:使用FlatMapElements

利用FlatMapElements将单个元素转换为元素流,实现拆分:

import org.apache.beam.sdk.Pipeline;
import org.apache.beam.sdk.io.TextIO;
import org.apache.beam.sdk.transforms.FlatMapElements;
import org.apache.beam.sdk.transforms.MapElements;
import org.apache.beam.sdk.values.PCollection;
import org.apache.beam.sdk.values.TypeDescriptors;
import java.util.Arrays;
import java.util.HashMap;
import java.util.Map;

public class SplitRecordExample {
    public static void main(String[] args) {
        Pipeline pipeline = Pipeline.create();
        
        PCollection<Map<String, String>> input = pipeline
            .apply("读取源数据", TextIO.read().from("input.txt"))
            .apply("解析为Map", MapElements.into(TypeDescriptors.maps(TypeDescriptors.strings(), TypeDescriptors.strings()))
                .via(line -> {
                    Map<String, String> record = new HashMap<>();
                    record.put("Column1", line.split(",")[0]); // 按需修改解析逻辑
                    return record;
                }));
        
        input.apply("拆分记录", FlatMapElements.into(TypeDescriptors.maps(TypeDescriptors.strings(), TypeDescriptors.strings()))
                .via(record -> {
                    if ("A".equals(record.get("Column1"))) {
                        Map<String, String> record1 = new HashMap<>(record);
                        record1.put("Column2", "A1");
                        Map<String, String> record2 = new HashMap<>(record);
                        record2.put("Column2", "A2");
                        return Arrays.asList(record1, record2);
                    } else {
                        return Arrays.asList(record);
                    }
                }))
            .apply("转换为文本", MapElements.into(TypeDescriptors.strings())
                .via(record -> record.get("Column1") + "," + record.getOrDefault("Column2", "")))
            .apply("输出结果", TextIO.write().to("output.txt"));
        
        pipeline.run().waitUntilFinish();
    }
}

方法2:自定义DoFn

继承DoFn,在processElement方法中多次调用output()输出记录:

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.MapElements;
import org.apache.beam.sdk.values.PCollection;
import org.apache.beam.sdk.values.TypeDescriptors;
import java.util.HashMap;
import java.util.Map;

public class SplitRecordDoFnExample {
    static class SplitRecordFn extends DoFn<Map<String, String>, Map<String, String>> {
        @ProcessElement
        public void processElement(ProcessContext c) {
            Map<String, String> record = c.element();
            if ("A".equals(record.get("Column1"))) {
                Map<String, String> record1 = new HashMap<>(record);
                record1.put("Column2", "A1");
                c.output(record1);
                
                Map<String, String> record2 = new HashMap<>(record);
                record2.put("Column2", "A2");
                c.output(record2);
            } else {
                c.output(record);
            }
        }
    }

    public static void main(String[] args) {
        Pipeline pipeline = Pipeline.create();
        
        PCollection<Map<String, String>> input = pipeline
            .apply("读取源数据", TextIO.read().from("input.txt"))
            .apply("解析为Map", MapElements.into(TypeDescriptors.maps(TypeDescriptors.strings(), TypeDescriptors.strings()))
                .via(line -> {
                    Map<String, String> record = new HashMap<>();
                    record.put("Column1", line.split(",")[0]);
                    return record;
                }));
        
        input.apply("拆分记录", ParDo.of(new SplitRecordFn()))
            .apply("转换为文本", MapElements.into(TypeDescriptors.strings())
                .via(record -> record.get("Column1") + "," + record.getOrDefault("Column2", "")))
            .apply("输出结果", TextIO.write().to("output.txt"));
        
        pipeline.run().waitUntilFinish();
    }
}

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.18 09:07:38