批处理流处理大数据【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam15/beam点击查看免费下载导读附加输出additional outputs也叫 tagged outputs 或 side outputs是 Apache Beam 中让一个ParDo变换同时产出多个PCollection的核心机制除了主输出main output之外你还可以声明任意多个带唯一标签的附加输出集合并在同一个DoFn的processElement/process方法内按条件把元素分发到不同集合。本文以 附加输出技术说明 为骨架结合仓库中的 Java SDK 与 Python SDK 源码、Kata 练习和 Tour of Beam 示例完整讲解从标签声明、输出发射到结果消费的完整链路。读完本文你将掌握何时该用附加输出来替代多个串联ParDo、如何在 Java 中通过TupleTagwithOutputTags声明多路输出、如何在 Python 中通过with_outputs()TaggedOutput实现同样的分支以及如何用PCollectionTuple/DoOutputsTuple分别取回各输出。一、什么是附加输出一个 ParDo 产生多个 PCollection在 Beam 的数据流模型中一次ParDo通常把一个输入PCollection的每个元素映射到一个输出元素。但很多真实场景需要在同一份输入上同时得到多种形态的结果例如把数值集合按阈值拆成小于等于 100和大于 100两组对每行文本同时产出分词结果主输出和字符数统计附加输出在一条流水线里既保留原始记录又产出清洗后的记录。附加输出正是为此设计的一个ParDo变换除了主输出PCollection外还可以产生任意数量的附加输出PCollection并把这些集合与主输出捆绑在一起返回。它本质上就是 Beam 实现**流水线分支pipeline branching**的机制——当需要把一个变换的输出拆分成多个PCollection或需要以不同格式产出结果时都优先考虑它。使用附加输出的一个显著收益是单次元素处理当每个元素的处理开销很大例如调用外部服务、进行复杂计算时用带附加输出的单个ParDo可以只遍历输入PCollection一次同时完成多种分发而不必为每种输出各写一个ParDo导致同一元素被重复处理多遍。产生附加输出有两个必要步骤为每个输出PCollection打上唯一标签tag在DoFn内部依据标签把元素发射到对应的输出集合。二、Java SDKTupleTag 声明输出MultiOutputReceiver 发射元素2.1 用 TupleTag 为每个输出打标签在 Java SDK 中附加输出通过 TupleTag 实现。TupleTagV是一个带泛型参数的、可序列化的类型化标签用来作为异质类型元组如PCollectionTuple的键其泛型参数让编译器能静态跟踪每个标签对应集合的元素类型。声明时有两个值得注意的约定见 TupleTag.java 的 javadoc输出标签作为 ParDo 输出时建议用匿名子类写法new TupleTagSomeType(){}这样标签类型信息能被类型描述符捕获便于 Beam 为输出结果推断默认 Coder输入标签用于取回结果用new TupleTag()即可额外实例化也无妨。从实现看TupleTag.java 的genId()会为每个标签生成唯一 id若标签在静态初始化块中创建如static final TupleTagT MY_TAG new TupleTag();则按类名序号分配确定性的 id否则用随机 nonce 拼上调用者信息降低跨 worker 序列化时的冲突概率。标签的相等性基于idequals/hashCode 实现因此可以在不同模块间仅凭 id 协调但同名标签必须携带相同的泛型类型否则会产生难以排查的运行时类型错误。2.2 用 withOutputTags 把标签绑定到 ParDo声明好标签后通过ParDo的.withOutputTags(...)方法把主输出标签和附加输出标签列表一起传入。附加输出标签使用 TupleTagList 组织可以链式.and(...)追加任意多个。withOutputTags返回的变换应用于输入PCollection后得到的是一个 PCollectionTuple ——一个异质类型的 PCollection 元组持有该ParDo的全部输出集合。下面是文档中的经典示例已修正标签类型并补全注释它实现了一个输入为字符串集合、输出为一个字符串主集合 一个字符串附加集合 一个整数附加集合的三路输出// 输入 PCollection元素为字符串。 PCollectionString input ...; // 主输出的标签字符串 PCollection。 final TupleTagString mainOutputTag new TupleTagString() {}; // 附加输出的标签字符串 PCollection。 final TupleTagString additionalOutputTagString new TupleTagString() {}; // 附加输出的标签整数 PCollection。 final TupleTagInteger additionalOutputTagIntegers new TupleTagInteger() {}; PCollectionTuple results input.apply(ParDo .of(new DoFnString, String() { // DoFn 主体继续在这里定义。 ... }) // 指定主输出的标签。 .withOutputTags(mainOutputTag, // 以 TupleTagList 指定两个附加输出的标签。 TupleTagList.of(additionalOutputTagString) .and(additionalOutputTagIntegers)));2.3 在 processElement 中向指定输出发射元素DoFn的processElement方法签名使用MultiOutputReceiver作为输出接收器注意发射多个输出时不能再使用OutputReceiverT单输出签名。通过out.get(tag).output(element)可以把元素发射到任意一个已声明的输出public void processElement(Element String word, MultiOutputReceiver out) { if (condition /* 主输出条件 */) { // 发射到主输出 out.get(mainOutputTag).output(word); } else { // 发射到附加的字符串输出 out.get(additionalOutputTagString).output(word); } if (condition /* 附加整数输出条件 */) { // 发射到附加的整数输出 out.get(additionalOutputTagIntegers).output(word.length()); } }2.4 用 PCollectionTuple.get(tag) 取回各路输出下游消费者只需用对应标签调用PCollectionTuple.get(tag)即可拿到独立的PCollection继续接变换。仓库中 Kata 练习Side Output 是一个可以直接运行的完整实现输入Create.of(10, 50, 120, 20, 200, 0)按 100 阈值拆成两路随后分别用Log.ofElements(...)打印TupleTagInteger numBelow100Tag new TupleTagInteger() {}; TupleTagInteger numAbove100Tag new TupleTagInteger() {}; PCollectionTuple outputTuple applyTransform(numbers, numBelow100Tag, numAbove100Tag); outputTuple.get(numBelow100Tag).apply(Log.ofElements(Number 100: )); outputTuple.get(numAbove100Tag).apply(Log.ofElements(Number 100: )); // applyTransform 内部 return numbers.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)));Tour of Beam 的 additional-outputs Java 示例 也采用完全相同的模式用自定义LogOutputTDoFn 输出日志可作为第二个对照实现参考。三、Python SDKwith_outputs() TaggedOutput 实现多路输出3.1 在 DoFn.process 中用 TaggedOutput 打标签发射Python SDK 中DoFn.process是一个生成器函数用yield发射元素。附加输出通过yield pvalue.TaggedOutput(tag, value)实现带上标签的 yield 会把元素送入对应标签的附加集合不带标签的普通yield value则进入主输出集合。文档中的SplitLinesToWordsFn是典型示例同时产出短词附加集合与字符数附加集合class SplitLinesToWordsFn(beam.DoFn): # 这些标签将用于标记此 DoFn 的各路输出。 OUTPUT_TAG_SHORT_WORDS tag_short_words OUTPUT_TAG_CHARACTER_COUNT tag_character_count def process(self, element): # 把字符数整数yield 到 OUTPUT_TAG_CHARACTER_COUNT 标签集合。 yield pvalue.TaggedOutput(self.OUTPUT_TAG_CHARACTER_COUNT, len(element)) words re.findall(r[A-Za-z\], element) for word in words: if len(word) 3: # 把 word yield 到 OUTPUT_TAG_SHORT_WORDS 标签集合。 yield pvalue.TaggedOutput(self.OUTPUT_TAG_SHORT_WORDS, word) else: # 普通 yield进入主集合。 yield word3.2 用 with_outputs() 声明标签并绑定 main 标签在ParDo变换上调用with_outputs()可以声明预期出现的标签。with_outputs 的签名与文档 说明如下*tags合法的附加输出标签列表如果声明了标签列表后续在流水线中使用未声明的标签会报错main...通过关键字参数指定主输出对应的标签名主输出本身没有真实标签这个字符串只是给它在DoOutputsTuple里取个名字返回值是 DoOutputsTuple 类型对象把所有输出集合捆绑在一起支持o.tag、o[tag]属性/索引访问也支持迭代所有标签校验规则main指定的标签不能同时出现在*tags中否则抛出ValueError见 core.py 的校验逻辑。文档中的使用示例with beam.Pipeline(optionspipeline_options) as p: lines p | ReadFromText(known_args.input) # with_outputs 允许访问 DoFn 的显式打标输出。 split_lines_result lines | beam.ParDo(SplitLinesToWordsFn()).with_outputs( SplitLinesToWordsFn.OUTPUT_TAG_SHORT_WORDS, SplitLinesToWordsFn.OUTPUT_TAG_CHARACTER_COUNT, mainwords, ) # split_lines_result 是 DoOutputsTuple 类型对象。 words, _, _ split_lines_result short_words split_lines_result[SplitLinesToWordsFn.OUTPUT_TAG_SHORT_WORDS] character_count split_lines_result.tag_character_count注意words, _, _ split_lines_result的解包顺序元组的第一个位置对应主输出mainwords其余位置按声明顺序对应附加输出而通过标签名访问result[tag]或result.tag_name则更直观且不受顺序影响。3.3 完整可运行的 Python 对照示例仓库中 Tour of Beam 的 additional-outputs Python 示例 提供了与 Java 版一一对应的完整流水线import apache_beam as beam from apache_beam import pvalue num_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 # 普通 yield进入主输出 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: )这里把mainnum_below_100_tag指定为主输出标签因此在DoOutputsTuple中可以用results[num_below_100]访问主集合用results[num_above_100]访问附加集合。Kata 的 Python Side Output 任务 提供了另一份练习版实现适合动手验证。四、附加输出 vs 其他分支手段如何选择Beam 中实现一个输入拆成多个输出还有其他手段理解差异有助于正确选型附加输出Side Outputs单个ParDo内部按元素条件分发输入元素只被处理一次适合处理开销大、需要按元素属性分流阈值、类型、异常分支的场景Partition变换按元素的分区函数把输入分成固定数量的子集合每个元素进入恰好一个分区适合按序数或可枚举类别均匀拆分它的每个元素也只处理一次Flatten与上述相反是把多个集合合并成一个不能用于分支多个独立的ParDo若每个分支都需要对原始输入做完全独立的变换且互斥可以串联多个ParDo但同一元素会被处理多次开销更大。附加输出真正的独特价值在于一份输入、一次遍历同时产出主输出与若干附加输出且各路输出的元素类型可以各不相同例如 Java 示例中字符串与整数并存这正是PCollectionTupleJava与DoOutputsTuplePython作为异质类型元组存在的意义。五、常用陷阱与最佳实践主输出也必须有标签Java 中.withOutputTags(mainTag, additionalTags)的第一个参数就是主输出标签它同样由TupleTag声明不要遗漏输出标签的泛型要精确附加输出元素类型与标签泛型不一致会在运行时暴露类型错误输入标签new TupleTag()与输出标签new TupleTagSomeType(){}的用法差异来自 TupleTag 的 javadoc 约定Python 中 main 标签不能与附加标签重复with_outputs(main...)与*tags出现同名会直接抛ValueErrorcore.py未声明的标签不要使用Python 中一旦在with_outputs()里声明了合法标签列表后续使用未声明的标签属于错误用法DoFn 输出接收器签名要匹配Java 中只要声明了附加输出processElement就应使用MultiOutputReceiver而非单输出的OutputReceiver否则无法按标签发射标签命名可读性Python 中标签直接决定DoOutputsTuple的属性名result.tag_character_count选用清晰、稳定的常量命名如类级常量能显著提升流水线可读性。六、深入阅读本文主题依据附加输出技术说明Java SDK 核心实现TupleTag、TupleTagList、PCollectionTupleMultiOutputReceiver定义于 DoFnPython SDK 核心实现with_outputs、pvalue.TaggedOutput / DoOutputsTuple可运行示例Java Kata Side Output、Python Kata Side Output、Tour of Beam Java 示例、Tour of Beam Python 示例赞分享批处理流处理大数据【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam15/beam点击查看免费下载相关推荐Apache Beam Go Katas使用 ParDo 附加输出Additional Outputs一次生成多条 PCollectionApache Beam Go Katas使用 ParDo 附加输出Additional Outputs一次生成多条 PCollection 导读 本篇围绕大数据批处理流处理数据工程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批处理流处理大数据上一篇抖音无水印下载终极指南3分钟掌握douyin-downloader完整使用教程下一篇Android 14下的自动化工具兼容性困局FGA启动异常的技术解析与实战修复创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考