首页/新闻资讯/正文详情

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

发布时间:2026/9/27 21:24:57 来源:云帆数科 栏目:资讯中心
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),仅供参考

相关推荐

LMDeploy INT4 模型量化和部署实战:AWQ/GPTQ 4bit 权重压缩、推理与服务
LMDeploy INT4 模型量化和部署实战:AWQ/GPTQ 4bit 权重压缩、推理与服务

人工智能大模型模型推理服务推理引擎本地部署模型量化 【免费下载链接】lmdeploy LMDeploy is a toolkit for compressing, deploying, and serving LLMs. 项目地址: https://gitcode.com/gh_mirrors/lm/lmdeploy 点击查看 免费下载 本文以 LMDeploy 的 w4a16&… · 2026/9/27 21:24:57

嵌入式AI的能效基准与标准化:每瓦Token数的工程含义
嵌入式AI的能效基准与标准化:每瓦Token数的工程含义

摘要:嵌入式AI正从“能跑”走向“好用”,但能效评估的标准化仍然缺失。elexcon2026将“AI端侧部署功耗控制”列为核心议题。端侧芯片的竞争核心正在从峰值算力转向“每瓦Token数”。本文从能效指标、测试方法和选型策略三个维度,分析嵌入式AI… · 2026/9/27 21:24:57

2026届六大AI辅助写作方案实测:TaoToken统一Key接入千笔AI、aipasspaper、豆包与DeepSeek的配置骨架
2026届六大AI辅助写作方案实测:TaoToken统一Key接入千笔AI、aipasspaper、豆包与DeepSeek的配置骨架

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views … · 2026/9/27 21:24:57

MCP 协议深入:用 TaoToken 统一 Key 打通 Anthropic Model Context Protocol 与 AI 工具链
MCP 协议深入:用 TaoToken 统一 Key 打通 Anthropic Model Context Protocol 与 AI 工具链

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views … · 2026/9/27 22:03:46

仁怀网站建设怎么选?避开备案坑的实战指南
仁怀网站建设怎么选?避开备案坑的实战指南

仁怀网站建设怎么选?避开备案坑的实战指南 备案流程一头雾水,导致网站上线周期从两周拖到两个月,这种痛点在仁怀本地企业中极其常见。很多老板想搞个官网展示酱酒品牌,结果卡在ICP备案环节,不知道材料怎么填,更不知道 怎么选… · 2026/9/27 22:03:46

【模型部署】用 OpenCV Python API 加载与运行 PyTorch 模型:TaoToken 统一 Key 配置与 ONNX 推理验证
【模型部署】用 OpenCV Python API 加载与运行 PyTorch 模型:TaoToken 统一 Key 配置与 ONNX 推理验证

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views … · 2026/9/27 22:03:46

5个来自谷歌的Agent Skill设计模式:用SKILL.md骨架配TaoToken跑通第一个Agent
5个来自谷歌的Agent Skill设计模式:用SKILL.md骨架配TaoToken跑通第一个Agent

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views … · 2026/9/27 22:03:39

Dumate 安装 superpowers-zh 汉化 skills:Trae 里用 npx 一键配置 TaoToken
Dumate 安装 superpowers-zh 汉化 skills:Trae 里用 npx 一键配置 TaoToken

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views … · 2026/9/27 22:03:39

企业网站首页应如何布局适合什么场景
企业网站首页应如何布局适合什么场景

5招搞定企业首页布局:不懂代码也能兼顾性能优化 很多老板盯着空白文档发呆,手里攥着预算却不敢动。不是不想做网站,是怕做出来的东西慢如蜗牛,或者被黑客挂了马。其实, 自己不会代码想做网站… · 2026/9/27 22:03:33

MATLAB雷达信号脉冲压缩仿真:LFM线性调频、匹配滤波与距离分辨率实现
MATLAB雷达信号脉冲压缩仿真:LFM线性调频、匹配滤波与距离分辨率实现

简介:这套Matlab仿真工具完整呈现雷达信号脉冲压缩过程,从线性调频(LFM)信号生成、目标回波仿真到匹配滤波压缩处理均有可运行代码支撑,面向电子信息工程、计算机、数学等专业学生,适用于课程设计、期末大作… · 2026/9/27 0:00:01

汕头网站建设制作厂家避坑指南:5大注意事项救急
汕头网站建设制作厂家避坑指南:5大注意事项救急

汕头网站建设制作厂家避坑指南:5大注意事项救急 改个需求建站公司拖一周,这种憋屈事我见得太多了。 很多汕头老板找本地建站团队,签合同前看着方案挺美,一上线就变脸。 今天不聊虚的,直接拆解找 汕头网站建设制作厂家 时的5个核心 注意事项… · 2026/9/27 0:00:01

多模态虚假新闻检测实战:BERT+ResNet双塔与对比学习
多模态虚假新闻检测实战:BERT+ResNet双塔与对比学习

简介:基于PyTorch的多模态虚假新闻检测项目完整代码包,面向自然语言处理与计算机视觉交叉方向的开发者、科研人员及毕业设计选题者,解决社交媒体中文本与图像联合识别虚假新闻的问题。系统以BERT预训练模型提取文本语义特征,以Res… · 2026/9/27 0:00:01

MATLAB雷达信号脉冲压缩仿真:LFM线性调频、匹配滤波与距离分辨率实现
MATLAB雷达信号脉冲压缩仿真:LFM线性调频、匹配滤波与距离分辨率实现

简介:这套Matlab仿真工具完整呈现雷达信号脉冲压缩过程,从线性调频(LFM)信号生成、目标回波仿真到匹配滤波压缩处理均有可运行代码支撑,面向电子信息工程、计算机、数学等专业学生,适用于课程设计、期末大作… · 2026/9/27 0:00:01

汕头网站建设制作厂家避坑指南:5大注意事项救急
汕头网站建设制作厂家避坑指南:5大注意事项救急

汕头网站建设制作厂家避坑指南:5大注意事项救急 改个需求建站公司拖一周,这种憋屈事我见得太多了。 很多汕头老板找本地建站团队,签合同前看着方案挺美,一上线就变脸。 今天不聊虚的,直接拆解找 汕头网站建设制作厂家 时的5个核心 注意事项… · 2026/9/27 0:00:01

多模态虚假新闻检测实战:BERT+ResNet双塔与对比学习
多模态虚假新闻检测实战:BERT+ResNet双塔与对比学习

简介:基于PyTorch的多模态虚假新闻检测项目完整代码包,面向自然语言处理与计算机视觉交叉方向的开发者、科研人员及毕业设计选题者,解决社交媒体中文本与图像联合识别虚假新闻的问题。系统以BERT预训练模型提取文本语义特征,以Res… · 2026/9/27 0:00:01

了解更多?预约专属演示

我们的顾问将为您一对一讲解产品与方案

企业微信二维码