如何在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
相关产品推荐
相关产品推荐

