ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

Apache Beam 核心转换 FlatMap(FlatMapElements)详解:一对多元素映射与实战

Apache Beam 核心转换 FlatMap(FlatMapElements)详解:一对多元素映射与实战 大数据批处理流处理数据工程【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam4/beam点击查看免费下载导读FlatMap是 Apache Beam 核心转换Core Transforms中处理一个输入元素产生零个或多个输出元素的标准方案。本篇技术指南以 Tour of Beam 交互式教程中flat-map-elements单元为骨架结合当前仓库中 Java 与 Python 双 SDK 的完整可运行示例及底层源码实现带你掌握FlatMap/FlatMapElements的写法、典型应用场景、与Map、ParDo的差异以及基于源码层面的执行原理。读完本文你将能独立写出句子分词按次数展开重复单词等一对多数据变换并理解其底层遍历发射机制。一、FlatMap 是什么与 Map 同源、但可一对多输出在原文档中FlatMap 被定义为It works likeMap elements, but inside the logic you can do complex operations like dividing the list into separate elements and processing即用法与 Map 相同但函数内部可以执行更复杂的操作例如将一个列表拆分成多个独立元素再逐个输出。核心区别可以这样理解MapMapElements/beam.Map对每个输入元素恰好输出一个输出元素一对一映射。相关讲解可参见同目录下的 map-elements 说明。FlatMapFlatMapElements/beam.FlatMap对每个输入元素用户函数返回一个可迭代集合Iterable该集合中的所有元素会被展平flatten后逐个写入输出PCollection一对多映射。从源码角度看Java 的 FlatMapElements.java 在类注释中明确了其定位mapping a simple function that returns iterables over the elements of a PCollection and merging the results——即对PCollection中每个元素应用一个返回 Iterable 的简单函数并将结果合并输出。其expand实现FlatMapElements.java#L158-L172本质上仍是基于ParDo的DoFn在processElement中调用映射函数得到IterableOutputT res随后通过for (OutputT output : res) { receiver.output(output); }逐项发射。这意味着FlatMap 是 ParDo 的一个便捷封装如果你需要完整的状态、定时器、侧输入等高级能力仍应直接使用ParDo可对照 pardo-one-to-many 说明 中手写DoFn实现同样效果的写法。二、双 SDK 基础示例将句子拆分成单词原文档给出了 Java 与 Python 两个 SDK 的核心示例。下面是在当前仓库中可完整运行的版本含main与日志输出分别位于 flat-map-elements/java-example/Task.java 与 flat-map-elements/python-example/task.py。Java使用FlatMapElements.into(...).via(...)import java.util.Arrays; import org.apache.beam.sdk.Pipeline; import org.apache.beam.sdk.options.PipelineOptions; import org.apache.beam.sdk.options.PipelineOptionsFactory; import org.apache.beam.sdk.transforms.*; import org.apache.beam.sdk.values.PCollection; import org.apache.beam.sdk.values.TypeDescriptors; public class Task { public static void main(String[] args) { PipelineOptions options PipelineOptionsFactory.fromArgs(args).create(); Pipeline pipeline Pipeline.create(options); // 输入两个句子 PCollectionString input pipeline.apply(Create.of(Apache Beam, Unified Batch and Streaming)); // applyTransform() 将 [sentences] 转换为 [words] PCollectionString output applyTransform(input); output.apply(Log, ParDo.of(new LogOutputString())); pipeline.run(); } // 核心转换把每个句子按空格拆成单词并展平输出 static PCollectionString applyTransform(PCollectionString input) { return input.apply( FlatMapElements.into(TypeDescriptors.strings()) .via(sentence - Arrays.asList(sentence.split( ))) ); } static class LogOutputT extends DoFnT, T { private String prefix; LogOutput() { this.prefix Processing element; } ProcessElement public void processElement(ProcessContext c) { LOG.info(prefix : {}, c.element()); } } }其中两个关键 API 的作用如下依据 FlatMapElements.java 源码注释FlatMapElements.into(TypeDescriptorOutputT outputType)L100-L103声明输出元素的类型描述符供 Beam 推断输出PCollection的 Coder。若省略into(...)则类型信息只能依赖函数本身推断因此建议始终显式指定输出类型。.via(ProcessFunctionInputT, ? extends IterableOutputT fn)L118-L121传入映射函数该函数必须返回一个 Iterable除ProcessFunction外还支持InferableFunction可携带输入/输出类型描述符见 L79-L88以及带上下文的Contextful版本可访问元素时间戳、窗口、侧输入等见 L129-L133。运行上述程序输入Apache Beam与Unified Batch and Streaming两个元素输出将是 7 个单词元素Apache、Beam、Unified、Batch、and、Streaming——这正是一个元素 → 多个元素的展平效果。Python使用beam.FlatMapimport apache_beam as beam with beam.Pipeline() as p: (p | beam.Create([Apache Beam, Unified Batch and Streaming]) # Lambda 函数返回句子的单词列表 | beam.FlatMap(lambda sentence: sentence.split()) | beam.Map(print))Python 侧beam.FlatMap的实现位于 sdks/python/apache_beam/transforms/core.py#L2064-L2079其文档注释明确约定传入的 callable 必须为输入PCollection中的每个元素返回一个 iterable这些 iterable 的元素会被展平进输出PCollection若未提供 callable使用默认identity则要求输入元素本身已是可迭代对象直接将其展平。此外FlatMap也接受*args/**kwargs并透传给用户函数因此可以传入额外参数实现参数化的转换逻辑。三、进阶实战其他类型与按次数展开模式原文档的 Playground 练习部分指出You can use other types instead ofInteger——即 FlatMap 的输入与输出类型完全由你决定不限于字符串。Java展开带计数的 KV 对原文档给出的进阶 Java 示例展示了输入类型为KVString, Integer单词与其出现次数、输出类型为String的场景将每个 KV 对按次数重复输出单词。PCollectionString splitWords input.apply( FlatMapElements.into(strings()).via((KVString, Integer wordWithCount) - { ListString words new ArrayList(); for (int i 0; i wordWithCount.getValue(); i) { words.add(wordWithCount.getKey()); } return words; }));这里via的 Lambda 参数类型为KVString, Integer返回ListString每个KV元素被映射成一个包含getValue()个重复单词的列表再被逐项展平输出。例如输入(Hello, 1), (World, 2), (How, 3), (are, 4), (you, 5)时输出将依次为Hello、World World、How How How……即每个单词出现其计数次。Python元组列表展开Python 侧等价写法如下with beam.Pipeline() as p: words_with_counts p | Create words with counts beam.Create([ (Hello, 1), (World, 2), (How, 3), (are, 4), (you, 5)]) split_words words_with_counts | Split words beam.FlatMap( lambda word_with_count: [word_with_count[0]] * word_with_count[1])beam.FlatMap的 lambda 接收一个(word, count)元组返回[word] * count列表从而把带计数的单词展平为重复的单词序列。这也是 WordCount 类经典流程分词 → 统计 → 展平中常用的还原技巧。四、FlatMap 与 Map、ParDo 的选型对照为帮助快速决策下表汇总了三者的定位差异依据本单元说明、map-elements 说明 与 pardo-one-to-many 说明转换输出数量典型写法适用场景ParDo自定义DoFn任意0..n可多路输出ParDo.of(new DoFn...{ ProcessElement ... out.output(...) })需要状态、定时器、侧输入、多输出标签等完整能力MapElements/beam.Map恰好 1 个MapElements.into(...).via(x - ...)/beam.Map(lambda x: ...)简单的一对一纯函数映射FlatMapElements/beam.FlatMap0..nIterable 展平FlatMapElements.into(...).via(x - Arrays.asList(...))/beam.FlatMap(lambda x: ...)一对多映射分词、展开、过滤后合并输出等需要特别说明的是Python 的beam.Map在实现上会把单值结果包装成单元素列表再走与FlatMap相同的展平路径因此二者在 Python 侧本质是同一机制的两个便捷入口而在 Java 侧FlatMapElements的expand中FlatMapElements.java#L135-L177会对函数返回值Iterable做循环output从而天然支持空列表 过滤掉该元素的语义——如果你需要基于某条件丢弃元素让函数返回空列表即可无需额外分支。五、深入源码FlatMapElements 的执行链路与异常处理执行链路从 FlatMapElements.java#L135-L177 的expand方法可以看出FlatMapElements的底层实现是校验fn非空checkArgument(fn ! null, .via() is required)即必须先调用via(...)指定映射函数否则直接报错将用户函数包装进一个内部FlatMapDoFn extends DoFnInputT, OutputT在ProcessElement中调用用户函数得到IterableOutputT res遍历res逐元素调用receiver.output(output)发射到输出PCollection。若传入的是Contextful函数需要访问上下文则在包装DoFn时还会通过withSideInputs(...)挂载其声明的侧输入FlatMapElements.java#L138-L157。另外FlatMapDoFn通过getInputTypeDescriptor()/getOutputTypeDescriptor()向 Beam 提供类型信息以推断 CoderFlatMapElements.java#L179-L202因此输出类型描述符into(...)缺失时在流水线构建期就可能报错。异常处理exceptionsInto/exceptionsVia除常规用法外Java SDK 的FlatMapElements还内置了失败处理扩展exceptionsInto(TypeDescriptor)与exceptionsVia(exceptionHandler)FlatMapElements.java#L219-L261会将映射过程中抛出的异常捕获并把ExceptionElementInputT异常实例 正在处理的输入元素交给用户提供的处理器输出到一个独立的失败PCollection。官方源码注释中的用法示例ResultPCollectionString, String result words.apply( FlatMapElements .into(TypeDescriptors.strings()) // 可能抛出 ArrayIndexOutOfBoundsException .via((String line) - Arrays.asList(Arrays.copyOfRange(line.split( ), 1, 5))) .exceptionsVia(new WithFailures.ExceptionAsMapHandlerString() {})); PCollectionString output result.output(); // 正常输出 PCollectionString failures result.failures(); // 失败元素集合其底层实现FlatMapElements.java#L310-L388使用ParDo.withOutputTags(...)将主输出与失败输出分离值得注意的是源码注释特别指出发射主输出必须放在try块之外以避免 runner 融合优化导致捕获到与本转换无关的异常。这是生产环境中处理脏数据的推荐模式。六、在 Playground 中运行与实验原文档指出该示例的完整代码位于 Playground 窗口中可直接运行并自由修改实验。对应到当前仓库单元元数据 unit-info.yaml 记录了该单元面向Java 与 Python两种 SDK复杂度标记为MEDIUM任务名为flat-map-elements完整可运行代码见 java-example/Task.java 与 python-example/task.py。建议的动手实验方向将sentence.split( )改为sentence.split()观察对同一句子的展平结果变化在 Java 版本中把返回类型改为Collections.emptyList()验证输入被过滤的语义将输入类型替换为整数集合例如Create.of(1, 2, 3)配合via(n - IntStream.range(0, n).boxed().collect(...))生成1, 2, 2, 3, 3, 3这样的三角序列尝试在 Python 版本中给beam.FlatMap传入带副参数的函数beam.FlatMap(lambda w, n: [w]*n, n3)体会*args/**kwargs透传能力。七、小结FlatMap/FlatMapElements是 Apache Beam 中处理一对多映射的标准工具它以极简的 Lambda 语法让开发者把元素拆分的意图直接表达在转换链上底层则由ParDo包装DoFn遍历发射 Iterable 中的每个元素Java 实现见 FlatMapElements.javaPython 实现见 core.py#L2064。掌握它你就能优雅地完成分词、展开、按条件过滤并输出等大量数据预处理场景当需求升级到需要上下文、侧输入或精细异常分流时还可以从Contextful、exceptionsVia等进阶 API 平滑过渡。进一步学习本单元属于 Tour of Beam 的core-transforms/map主题同主题下还有 map-elements、pardo-one-to-one、pardo-one-to-many、group-by-key 与 co-group-by-key 等单元可对照学习元素映射与分组体系的完整脉络。赞分享大数据批处理流处理数据工程【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam4/beam点击查看免费下载相关推荐Apache Beam Kotlin Katas 实战用 FlatMapElements 实现一对多元素映射Apache Beam Kotlin Katas 实战用 FlatMapElements 实现一对多元素映射 本文以 Apache Beam 官方 Kotli大数据批处理流处理数据工程Apache Beam Java Kata 实战用 FlatMapElements 实现一对多one-to-many映射Apache Beam Java Kata 实战用 FlatMapElements 实现一对多one to many映射 FlatMapElements大数据批处理流处理数据工程Apache Beam Python FlatMap 实战用一对多映射与轻量 DoFn 简化数据处理Apache Beam Python FlatMap 实战用一对多映射与轻量 DoFn 简化数据处理 Apache Beam 的 Python SDK 提供了大数据批处理流处理数据工程上一篇5步掌握Sketch MeaXure彻底解决设计开发沟通障碍下一篇终极Sunshine游戏串流卸载指南如何彻底清理并释放系统资源创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表