ARTICLE DETAIL

资讯详情

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

Apache Beam Kotlin 实战:使用 ParDo 与 MultiOutputReceiver 实现多输出(Side Output)Kata 全解

Apache Beam Kotlin 实战:使用 ParDo 与 MultiOutputReceiver 实现多输出(Side Output)Kata 全解 大数据批处理流处理数据工程【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam4/beam点击查看免费下载Apache Beam 的统一编程模型中ParDo是进行逐元素并行处理的核心转换。本指南围绕 Beam KatasKotlin 版Side Output 这一课讲解如何让一个ParDo在产出主输出 PCollection 的同时通过TupleTag与MultiOutputReceiver产生任意数量的附加输出并完成把大于 100 的数字分流到独立 PCollection的 Kata 任务。读完本文你将掌握多输出ParDo的完整写法、PCollectionTuple的取用方式以及对应的测试验证方法。一、任务背景Katas 中的 Side Output 一课本 Kata 位于仓库的 Kotlin 学习路径中任务说明learning/katas/kotlin/Core Transforms/Side Output/Side Output/src/task.md待补全的代码骨架Task.kt隐藏的单元测试TaskTest.kt课程元数据task-info.yamltype: edu标记了TODO()占位符位置任务说明原文指出核心概念ParDo 总是产生一个主输出 PCollection作为apply的返回值但也可以让这个 ParDo 产生任意数量的附加输出 PCollection。如果选择多输出你的 ParDo 会将所有输出 PCollection包括主输出打包在一起返回。Kata 目标非常明确为你的 ParDo 实现一个附加输出用于接收大于 100 的数字。注意Side Output多输出与另一课 Side Input旁路输入把 PCollection 作为额外输入传入 DoFn是两种不同机制不要混淆。二、核心概念与 API 拆解2.1 主输出与附加输出在 Beam 中普通ParDo的返回值就是主输出PCollectionOutputT。一旦通过.withOutputTags(...)声明了附加输出返回值就变成PCollectionTuple——一个按TupleTag索引、捆绑了所有输出的容器见 ParDo.java 的类注释可选地一个 ParDo 转换可以产生多个输出 PCollection包括一个主输出PCollectionOutputT以及任意数量的附加输出 PCollection每个附加输出用一个不同的TupleTag标识并捆绑在一个PCollectionTuple中。附加输出所需的 TupleTag 通过调用SingleOutput#withOutputTags指定。2.2 TupleTag输出的身份证每个输出都必须绑定一个TupleTagT它是输出元素的类型标签同时承担类型信息与运行时标识的双重职责。Kata 中定义了两个标签val numBelow100Tag object : TupleTagInt() {} val numAbove100Tag object : TupleTagInt() {}这里有一个关键细节来自 ParDo.java 源码注释输出用的 TupleTag 必须实例化为匿名子类尾部带{}。原因在于 Beam 需要从TupleTagT的泛型参数推断附加输出 PCollection 的 Coder匿名子类会阻断 Java/Kotlin 的泛型类型推断从而强制显式写出类型参数保证 Coder 推断成功。若直接使用TupleTagInt()这种非匿名形式运行时将拿不到完整的类型信息。2.3 MultiOutputReceiver按标签发射元素在ProcessElement方法中多输出模式需要额外注入一个MultiOutputReceiver参数。其接口定义位于 DoFn.javapublic interface MultiOutputReceiver { T OutputReceiverT get(TupleTagT tag); T OutputReceiverRow getRowReceiver(TupleTagT tag); }get(tag)返回指定标签对应的OutputReceiverT通过它调用.output(element)把元素发射到该输出getRowReceiver(tag)当目标输出注册了 SchemaRow时使用。因此DoFnInt, Int的主输出类型参数第二个Int依然存在但多输出场景下你实际上通过MultiOutputReceiver按标签发射而不是使用context.output(...)直接写主输出。三、Kata 完整解法与逐行解析3.1 完整可运行代码applyTransform的完整实现即 TODO 占位符处应补全的内容Task.ktfun applyTransform( numbers: PCollectionInt, numBelow100Tag: TupleTagInt, numAbove100Tag: TupleTagInt ): PCollectionTuple { return numbers.apply(ParDo.of(object : DoFnInt, Int() { ProcessElement fun processElement(context: ProcessContext, out: MultiOutputReceiver) { val number context.element() if (number 100) { out.get(numBelow100Tag).output(number) } else { out.get(numAbove100Tag).output(number) } } }).withOutputTags(numBelow100Tag, TupleTagList.of(numAbove100Tag))) }3.2 关键点拆解1withOutputTags声明输出集合.withOutputTags(mainOutputTag, TupleTagList.of(additionalOutputTag))用于声明本次多输出的标签集合。源码 ParDo.java 中它的签名与行为如下public MultiOutputInputT, OutputT withOutputTags( TupleTagOutputT mainOutputTag, TupleTagList additionalOutputTags) { return new MultiOutput(fn, sideInputs, mainOutputTag, additionalOutputTags, fnDisplayData); }第一个参数是主输出标签本 Kata 中为numBelow100Tag第二个参数是附加输出标签列表通过TupleTagList.of(tag).and(tag)...可链式追加任意多个标签。2按条件分流ProcessElement中每个元素依据阈值 100 走不同分支number 100→ 通过out.get(numBelow100Tag).output(number)发往主输出number 100→ 通过out.get(numAbove100Tag).output(number)发往附加输出。3PCollectionTuple 按标签取流在main中applyTransform返回的PCollectionTuple通过.get(tag)取出各分支 PCollection再分别接上日志打印转换Log.ktoutputTuple.get(numBelow100Tag).apply(Log.ofElements(Number 100: )) outputTuple.get(numAbove100Tag).apply(Log.ofElements(Number 100: ))输入集合为Create.of(10, 50, 120, 20, 200, 0)Task.kt 第 37 行运行后控制台将分别打印Number 100: 0 Number 100: 10 Number 100: 20 Number 100: 50 Number 100: 120 Number 100: 2003.3 输出 Coders 的注意事项从 ParDo.java 源码 可以确认两条 Coder 推断规则主输出的 Coder 从DoFnInputT, OutputT的具体类型推断每个附加输出的 Coder 从对应TupleTagAdditionalOutputT的具体类型推断这就要求 TupleTag 必须写成匿名子类形式见 2.2 节。本 Kata 中主输出与附加输出的元素类型都是IntCoder 均能顺利推断为VarIntCoder。四、单元测试验证分流逻辑TaskTest.kt测试文件用 Beam 的测试框架验证了分流正确性Test fun core_transforms_side_output_side_output() { val numbers testPipeline.apply(Create.of(10, 50, 120, 20, 200, 0)) val numBelow100Tag object : TupleTagInt() {} val numAbove100Tag object : TupleTagInt() {} val resultsTuple applyTransform(numbers, numBelow100Tag, numAbove100Tag) PAssert.that(resultsTuple.get(numBelow100Tag)).containsInAnyOrder(0, 10, 20, 50) PAssert.that(resultsTuple.get(numAbove100Tag)).containsInAnyOrder(120, 200) testPipeline.run().waitUntilFinish() }测试要点使用TestPipelineRule构建测试管线构造与main相同的输入(10, 50, 120, 20, 200, 0)通过PAssert.that(...).containsInAnyOrder(...)断言每个输出流的内容——主输出必须恰好包含{0, 10, 20, 50}附加输出必须恰好包含{120, 200}注意PAssert断言的是无序集合即使ParDo内元素处理顺序不确定断言依然稳定成立。这组测试同时也反向证明了MultiOutputReceiverwithOutputTags的调用契约只有标签集合声明完整、发射路径与标签一一对应两个分支的输出才能被正确分离与取出。五、Katas 教学机制与多输出扩展5.1 教学机制task-info.yaml表明这是一个edu教育类型任务Task.kt中长度为 423 字节的TODO()占位符是学员需要补全的区域测试文件默认对学员隐藏。同一小节还包含 DoFn Additional Parameters、Side Input 等相关课程共同构成 Core Transforms 的完整练习链。对应地Java 版 Task.java 提供了等价的MultiOutputReceiver用法可作为跨语言对照。5.2 多输出不止两个分支TupleTagList支持链式追加任意数量的标签。参考 ParDo.java 的官方示例一个DoFn甚至可以同时发射主输出、多个附加输出甚至存在声明了但无人消费的输出——源码明确指出未消费的输出无须显式列出。典型应用场景包括按业务规则把数据拆分为正常数据流与异常/告警数据流分别走不同的下游处理在同一个 DoFn 内同时产出处理结果与处理指标如元素计数将无法解析的脏数据单独引到旁路流供后续审计或修复避免污染主链路。5.3 与 Side Input 的区分同小节的 Side Input 课程解决的是把 PCollection 作为附加输入读入 DoFn的问题而本课的 Side Output 解决的是从 DoFn 额外输出多个 PCollection的问题。二者是 Beam 数据流动的进与出两个方向常可组合使用例如把旁路输入的分组结果作为侧输入在主 DoFn 中结合多输出完成复杂分流。六、常见错误与排查建议TupleTag 未写成匿名子类TupleTagInt()直接实例化会导致附加输出 Coder 推断失败或类型信息丢失必须写成object : TupleTagInt() {}Java 中为new TupleTagInteger() {}。忘记注入 MultiOutputReceiver 参数多输出 DoFn 的ProcessElement必须声明out: MultiOutputReceiver参数否则无法按标签发射元素。发射到了未声明的标签withOutputTags未列出的标签即使被发射也不会出现在返回的PCollectionTuple中务必保证标签集合声明完整。把主输出标签当作普通标签withOutputTags的第一个参数是主输出标签其余标签放入TupleTagList标签顺序决定PCollectionTuple的结构读取时始终通过tuple.get(tag)而非位置索引。结语通过本 Kata你完成了 Beam Kotlin 多输出ParDo的完整闭环从TupleTag定义、MultiOutputReceiver按标签发射到withOutputTags声明标签集合、PCollectionTuple取流再到PAssert单元测试验证。这套一进多出的模式是 Beam 生产管线中数据分流、旁路告警、多路复用的基石建议在此基础上继续练习 Composite Transform 与 Partition 等分支类课程构建更完整的 Core Transforms 能力图谱。赞分享大数据批处理流处理数据工程【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam4/beam点击查看免费下载相关推荐Apache Beam Kotlin Katas 实战用 ParDo 的 Side Output额外输出实现数据分流Apache Beam Kotlin Katas 实战用 ParDo 的 Side Output额外输出实现数据分流 本指南以 Apache Beam 仓大数据批处理流处理数据工程Apache Beam Java 实战用 ParDo 多路输出Side Output拆分大于 100 的数字流Apache Beam Java 实战用 ParDo 多路输出Side Output拆分大于 100 的数字流 Apache Beam 的 ParDo 变大数据批处理流处理数据工程使用 ParDo 与 DoFn 在 Apache Beam Kotlin 中实现过滤转换Kata 实战指南使用 ParDo 与 DoFn 在 Apache Beam Kotlin 中实现过滤转换Kata 实战指南 本篇技术指南围绕 Apache Beam 官方 K大数据批处理流处理数据工程上一篇MuJoCo中物体打滑的完整止滑调参指南下一篇Cherry Studio 中的 Claude Code MCP 服务器选型与配置指南从推荐清单到运行时实现创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表