ARTICLE · INTELLIGENCE

战地情报 · 详情页

来自尧图项目组的一线实战观察与深度解析

Apache Beam 多输出(Additional Outputs)实战指南:Go / Java / Python 三 SDK 的 ParDo 多 PCollection 输出机制详解

Apache Beam 多输出(Additional Outputs)实战指南:Go / Java / Python 三 SDK 的 ParDo 多 PCollection 输出机制详解 批处理流处理大数据【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam15/beam点击查看免费下载ParDo是 Apache Beam 中最核心的按元素处理变换而多输出Additional Outputs让单个ParDo能够一次性产生多个不同类型、不同语义的PCollection避免为每个分支重复扫描数据。本篇以 Apache Beam 仓库中 core-transforms/additional-outputs 学习单元为骨架完整讲解 Go、Java、Python 三种 SDK 实现多输出的不同机制Go 的位置式 Emitter、Java 的TupleTagPCollectionTuple、Python 的with_outputsTaggedOutput并深入到 pardo.go、pvalue.py 等源码印证底层实现最后给出可运行的 Playground 完整练习代码。读完本文你将掌握在任意一种 SDK 中按条件把一条数据流拆分成多条输出流、以及如何在DoFn中访问时间戳、窗口与 Pane 信息等附加参数的完整能力。为什么需要多输出一次遍历多条数据流在数据处理中一个非常常见的需求是根据某种条件把输入数据分流。例如把单词按长度分成短词与长词、把整数按是否超过阈值分成两个集合、把日志按级别拆成多条下游路径。ParDo总是会产出主输出main output但 Beam 允许同一个DoFn额外产出任意数量的附加输出PCollection且这些输出可以是不同的数据类型。相比把条件判断搬到下游重复计算、重复读取数据多输出让一次元素处理就能把元素分发到多个目的地这是处理大数据或需要将数据拆分到不同集合时的推荐做法。该学习单元在仓库中的元数据 unit-info.yaml 标记为complexity: MEDIUM并同时支持 Java、Python、Go 三种 SDK。三种 SDK 的多输出 API 形态差异较大下面分别展开。Go SDK位置式多输出Positional OutputsParDo2 ~ ParDo7 与 ParDoN按输出数量选择函数Go SDK 的DoFn总能产生一个输出PCollection而通过选择合适的ParDo变体一个DoFn可以产生任意数量的附加输出PCollection甚至一个都不产生beam.ParDo2两个输出PCollectionbeam.ParDo3三个输出依此类推直到beam.ParDo7需要更多时使用beam.ParDoN它返回[]beam.PCollection切片一个输出都不需要时使用beam.ParDo0。这些变体在 sdks/go/pkg/beam/pardo.go 中均有明确实现ParDoN在第 129 行ParDo0在第 134 行ParDo2~ParDo7分别在第 439 ~ 484 行。它们都基于内部的TryParDo展开最终返回[]PCollection见 pardo.go 的返回收集逻辑ParDo2等只是把切片解包成多个返回值// 来自 sdks/go/pkg/beam/pardo.go 的函数签名 func ParDoN(s Scope, dofn any, col PCollection, opts ...Option) []PCollection func ParDo0(s Scope, dofn any, col PCollection, opts ...Option) func ParDo2(s Scope, dofn any, col PCollection, opts ...Option) (PCollection, PCollection) func ParDo3(s Scope, dofn any, col PCollection, opts ...Option) (PCollection, PCollection, PCollection)TagsGo 不用标签用位置顺序Go SDK 的多个输出不使用命名标签而是依赖位置顺序ParDo返回的PCollection顺序与DoFn中 emit 函数参数的出现顺序一一对应。下面的示例中processWords的三个 emitter 按参数位置对应beam.ParDo3返回的三个PCollection// beam.ParDo3 returns PCollections in the same order as // the emit function parameters in processWords. // Now below has the output of s, above has the output of processWords, and marked has the output of words below, above, marked : beam.ParDo3(s, processWords, words) // processWordsMixed uses both a standard return and an emitter function. // The standard return produces the first PCollection from beam.ParDo2, // and the emitter produces the second PCollection. length, mixedMarked : beam.ParDo2(s, processWordsMixed, words)注意第二个示例体现了一条重要规则DoFn的标准返回值永远对应beam.ParDo返回的第一个PCollection其他 emitter 则按它们在DoFn方法参数中的定义顺序输出到各自的PCollection。在 DoFn 中向多个输出发射元素Go 的DoFn通过发射函数emitter向指定输出写元素按需调用 emitter 即可向其对应的PCollection产生 0 个或多个元素同一个值可以被多个 emitter 同时发射与单输出一致发射后不应再修改该值。// processWords is a DoFn that has 3 output PCollections. The emitter functions // are matched in positional order to the PCollections returned by beam.ParDo3. func processWords(word string, emitBelowCutoff, emitAboveCutoff, emitMarked func(string)) { const cutOff 5 if len(word) cutOff { emitBelowCutoff(word) } else { emitAboveCutoff(word) } if isMarkedWord(word) { emitMarked(word) } }所有 emitter 都应通过泛型的register.EmitterX[...]函数注册这样会优化 emitter 的运行时执行效率这也是 Beam Go SDK 推荐的性能实践。DoFn 中的附加参数Context、时间戳、窗口与 Pane除了元素本身Beam 会自动为DoFn的ProcessElement方法注入其他参数任何组合都可以按标准顺序添加context.Context用于支持统一日志和用户自定义指标。按 Go 惯例如果存在必须是DoFn方法的第一个参数func MyDoFn(ctx context.Context, word string) string { ... }Timestamp时间戳要访问输入元素的事件时间戳在元素参数之前添加beam.EventTime参数func MyDoFn(ts beam.EventTime, word string) string { ... }Window窗口要访问元素落入的窗口在元素参数之前添加beam.Window参数。注意如果一个元素同时落入多个窗口例如使用SlidingWindows滑动窗口时ProcessElement会对该元素按窗口各调用一次。由于beam.Window是接口可以类型断言到具体窗口实现例如固定窗口场景下断言为window.IntervalWindowfunc MyDoFn(w beam.Window, word string) string { iw : w.(window.IntervalWindow) ... }PaneInfoPane 信息使用触发器时Beam 提供beam.PaneInfo对象包含当前这次触发的信息。借助它可判断本次触发是早期early还是晚期late触发以及该窗口针对该 key 已经触发了多少次func extractWordsFn(pn beam.PaneInfo, line string, emitWords func(string)) { if pn.Timing typex.PaneEarly || pn.Timing typex.PaneOnTime { // ... perform operation ... } if pn.Timing typex.PaneLate { // ... perform operation ... } if pn.IsFirst { // ... perform operation ... } if pn.IsLast { // ... perform operation ... } words : strings.Split(line, ) for _, w : range words { emitWords(w) } }Java SDKTupleTag PCollectionTuple 的带键多输出核心概念PCollectionTuple 与 TupleTagJava 的ParDo同样始终输出主输出PCollectionapply的返回值但可以强制它输出任意数量的附加PCollection。当选择多输出时ParDo会把所有输出包括主输出合并在一起返回得到PCollectionTuple再通过TupleTag取出其中任意一个PCollection。PCollectionTuple异构类型PCollection的不可变元组以TupleTag作为键。它既可作为接收/创建多个不同类型输入或输出的PTransform的输入或输出典型场景就是带多输出的ParDo。TupleTagT异构类型元组如PCollectionTuple的带类型标签其泛型参数用于追踪存储在元组中数据的静态类型。用 MultiOutputReceiver 向多个输出发射在DoFn中通过MultiOutputReceiver out与out.get(tag).output(value)将元素发射到指定输出ParDo.of(new DoFnString, String() { public void processElement(Element String word, MultiOutputReceiver out) { if (word.length() wordLengthCutOff) { // Emit short word to the main output. // In this example, it is the output with tag wordsBelowCutOffTag. out.get(wordsBelowCutOffTag).output(word); } else { // Emit long word length to the output with tag wordLengthsAboveCutOffTag. out.get(wordLengthsAboveCutOffTag).output(word.length()); } if (word.startsWith(MARKER)) { // Emit word to the output with tag markedWordsTag. out.get(markedWordsTag).output(word); } }}) .withOutputTags(wordsBelowCutOffTag, // Specify the tags for the two additional outputs as a TupleTagList. TupleTagList.of(wordLengthsAboveCutOffTag).and(markedWordsTag));关键点withOutputTags的第一个参数是主输出标签其余附加输出标签通过TupleTagList.of(...).and(...)链式指定上述示例还展示了多输出的类型灵活性主输出是String而wordLengthsAboveCutOffTag对应的输出类型是Integer。Python SDKwith_outputs TaggedOutput声明输出标签with_outputs 返回 DoOutputsTuplePython 中要让ParDo产生多个输出需在ParDo上调用with_outputs()并声明期望的输出标签。with_outputs()返回一个DoOutputsTuple对象源码定义见 sdks/python/apache_beam/pvalue.py标签会作为该对象上的属性也可以像字典一样按下标索引访问对应的输出PCollectionresults ( input | beam.ParDo(ProcessWords(), cutoff_length2, markerx).with_outputs(above_cutoff_lengths,marked strings,mainbelow_cutoff_strings)) below results.below_cutoff_strings above results.above_cutoff_lengths marked results[marked strings] # indexing works as wellDoOutputsTuple同时也是可迭代对象迭代顺序与传入with_outputs()的标签顺序一致若指定了主标签main则主输出排在最前below, above, marked (input | beam.ParDo(ProcessWords(), cutoff_length2, markerx) .with_outputs(above_cutoff_lengths,marked strings,mainbelow_cutoff_strings))关于with_outputs()的约束其实现位于 sdks/python/apache_beam/transforms/core.py若声明了有效标签列表则后续在管道中使用未声明标签会报错且主输出标签main不能与附加输出标签重复否则抛出ValueError。在 DoFn 中发射pvalue.TaggedOutput在DoFn内部通过把值输出标签包装进pvalue.TaggedOutput来把元素发射到指定输出class ProcessWords(beam.DoFn): def process(self, element, cutoff_length, marker): if len(element) cutoff_length: # Emit this short word to the main output. yield element else: # Emit this words long length to the above_cutoff_lengths output. yield pvalue.TaggedOutput(above_cutoff_lengths, len(element)) if element.startswith(marker): # Emit this word to a different output with the marked strings tag. yield pvalue.TaggedOutput(marked strings, element)TaggedOutput的实际实现位于 sdks/python/apache_beam/pvalue.py。值得强调的是原文档正文中提到的pvalue.OutputValue包装类在当前仓库的 pvalue.py 中实际对应的是pvalue.TaggedOutput两者同属pvalue模块的输出包装机制以仓库源码中实际可用的TaggedOutput为准。Map / FlatMap 同样支持多输出多输出不局限于ParDoMap和FlatMap也可以。而且用FlatMap时标签无需提前声明def even_odd(x): yield pvalue.TaggedOutput(odd if x % 2 else even, x) if x % 10 0: yield x results input | beam.FlatMap(even_odd).with_outputs() evens results.even odds results.odd tens results[None] # the undeclared main output这个示例同时展示了三个细节动态决定输出标签、未声明的标签直接可用、以及用results[None]访问未声明的主输出。Playground 练习整数与字符串的双路分流该学习单元的练习题要求applyTransform()接收一个整数列表输出两个PCollection——一个存放大于 100 的数一个存放小于等于 100 的数进阶版则改为对字符串按大小写分流。完整的可运行代码分别位于 go-example/main.go、java-example/Task.java 和 python-example/task.py可直接在 Tour of Beam 的 Playground 窗口中运行与修改。Go 版本ParDo2 双 Emitter// 主程序构造输入并消费两个输出 p, s : beam.NewPipelineWithRoot() input : beam.Create(s, 10, 50, 120, 20, 200, 0) numBelow100, numAbove100 : applyTransform(s, input) debug.Printf(s, Number 100: %v, numBelow100) debug.Printf(s, Number 100: %v, numAbove100) // 核心一个 DoFn 两个 emitter按位置对应两个输出 func applyTransform(s beam.Scope, input beam.PCollection) (beam.PCollection, beam.PCollection) { return beam.ParDo2(s, func(element int, numBelow100, numAbove100 func(int)) { if element 100 { numBelow100(element) return } numAbove100(element) }, input) }练习的字符串进阶版需要先用ParDo把句子拆成单词用strings.Split并按空格过滤再用ParDo2判断element strings.Title(element)即是否首字母大写决定发射到大写或小写输出。Java 版本withOutputTags MultiOutputReceiverTupleTagInteger numBelow100Tag new TupleTagInteger() {}; TupleTagInteger numAbove100Tag new TupleTagInteger() {}; PCollectionTuple outputTuple applyTransform(input, numBelow100Tag, numAbove100Tag); outputTuple.get(numBelow100Tag).apply(Log Number 100: , ParDo.of(new LogOutputInteger(Number 100: ))); outputTuple.get(numAbove100Tag).apply(Log Number 100: , ParDo.of(new LogOutputInteger(Number 100: ))); static PCollectionTuple applyTransform( PCollectionInteger input, TupleTagInteger numBelow100Tag, TupleTagInteger numAbove100Tag) { return input.apply(ParDo.of(new DoFnInteger, Integer() { ProcessElement public void processElement(Element Integer number, MultiOutputReceiver out) { if (number 100) { out.get(numBelow100Tag).output(number); } else { out.get(numAbove100Tag).output(number); } } }).withOutputTags(numBelow100Tag, TupleTagList.of(numAbove100Tag))); }字符串版本把输入换成句子先用FlatMapElements.into(TypeDescriptors.strings())配合正则[^\p{L}]分词再在DoFn中用element.equals(element.toLowerCase())判断是否全小写分别发射到lowerCase/upperCase两个标签。Python 版本TaggedOutput with_outputsnum_below_100_tag num_below_100 num_above_100_tag num_above_100 class ProcessNumbersDoFn(beam.DoFn): def process(self, element): if element 100: yield element # 主输出num_below_100 else: yield pvalue.TaggedOutput(num_above_100_tag, element) with beam.Pipeline() as p: results (p | beam.Create([10, 50, 120, 20, 200, 0]) | beam.ParDo(ProcessNumbersDoFn()).with_outputs(num_above_100_tag, mainnum_below_100_tag)) results[num_below_100_tag] | Log nums below 100 Output(prefixnum_below_100: ) results[num_above_100_tag] | Log nums above 100 Output(prefixnum_above_100: )注意 Python 版通过mainnum_below_100_tag显式给主输出命名之后就能用results[num_below_100_tag]与results[num_above_100_tag]分别访问两条输出流。多输出的关键规则与常见陷阱结合三种 SDK 的文档说明与源码使用多输出时值得牢记以下要点主输出始终存在且固定Go 中DoFn的标准返回值固定对应第一个返回的PCollectionPython 中主输出默认不带标签用None索引访问可通过main命名Java 中withOutputTags的第一个参数即主输出标签。输出类型可以异构PCollectionTuple、[]PCollection、DoOutputsTuple都允许每个输出持有不同类型的数据Java 的TupleTagT泛型即为追踪类型而设计。同一个值可以发射到多个输出Go 与 Python 的示例都展示了同一元素同时进入两个输出的写法如既符合长度条件又带 MARKER 前缀。发射后不可再修改元素Go 文档明确强调从任意 emitter 发射之后不要再修改该值。Go 的 emitter 需用register.EmitterX[...]注册以优化运行时执行。Python 声明标签后的约束with_outputs若声明了标签列表后续使用未声明标签属于错误但FlatMap场景下不预先声明标签也可以动态产生新输出。附加参数注入GoDoFn可自由组合注入context.Context、beam.EventTime、beam.Window、beam.PaneInfo其中Context必须是首个参数Window可类型断言到具体窗口类型滑动窗口下同一元素会被每个窗口各调用一次ProcessElement。延伸阅读学习单元文档learning/tour-of-beam/learning-content/core-transforms/additional-outputs/description.md单元元数据SDK 支持与复杂度learning/tour-of-beam/learning-content/core-transforms/additional-outputs/unit-info.yamlGo 可运行示例learning/tour-of-beam/learning-content/core-transforms/additional-outputs/go-example/main.goJava 可运行示例learning/tour-of-beam/learning-content/core-transforms/additional-outputs/java-example/Task.javaPython 可运行示例learning/tour-of-beam/learning-content/core-transforms/additional-outputs/python-example/task.pyGo SDK 的 ParDo 变体源码sdks/go/pkg/beam/pardo.goPython SDK 的DoOutputsTuple与TaggedOutput源码sdks/python/apache_beam/pvalue.pyPython SDK 的with_outputs源码sdks/python/apache_beam/transforms/core.py赞分享批处理流处理大数据【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam15/beam点击查看免费下载相关推荐Apache Beam 多输出Additional Outputs实战指南Go / Java / Python 三 SDK 的 ParDo 多 PCollection 输出Apache Beam 多输出Additional Outputs实战指南Go / Java / Python 三 SDK 的 ParDo 多 PColl大数据批处理流处理数据工程Apache Beam Go SDK 实战使用 ParDo 额外输出Additional Outputs分流多条 PCollectionApache Beam Go SDK 实战使用 ParDo 额外输出Additional Outputs分流多条 PCollection 本文基于 Apa批处理流处理大数据Apache Beam Go Katas使用 ParDo 附加输出Additional Outputs一次生成多条 PCollectionApache Beam Go Katas使用 ParDo 附加输出Additional Outputs一次生成多条 PCollection 导读 本篇围绕大数据批处理流处理数据工程上一篇AMD ROCm完整指南如何解决AI框架GPU识别问题并优化性能下一篇告别复杂集成use-mcp让React应用轻松接入AI能力创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
RELATED READING

延伸阅读

更多一线实战笔记与深度复盘,助您持续精进