ARTICLE DETAIL

资讯详情

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

Apache Beam Go SDK 实战:使用 ParDo2 实现 Additional Outputs 多路输出

Apache Beam Go SDK 实战:使用 ParDo2 实现 Additional Outputs 多路输出 【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam18/beam点击查看免费下载导读本文基于 Apache Beam Go SDK 的 Code Katas 教程中 Additional Outputs额外输出一课系统讲解如何让一个ParDo变换同时产生多个PCollection输出。通过一个把大于 100 的数字单独输出到另一条 PCollection的实战练习你将掌握beam.ParDo2的 DoFn 签名规范、完整可运行的管道代码、测试验证方法并深入理解 Go SDK 中ParDo系列变换ParDo1到ParDo7的底层实现原理。一、什么是 Additional Outputs在之前的课程中我们通常把一个DoFn应用到ParDo上输出一个单一的PCollection。但实际上ParDo变换可以产生零个或多个输出 PCollection——这就是 Beam 的 Additional Outputs额外输出能力。在真实的数据处理场景中我们经常需要把输入数据按某种业务规则分流例如把日志按严重级别拆分、把交易按金额区间拆分、把数据按是否合法拆分。如果为每种情况分别写一个ParDo遍历输入会造成多次读取同一份数据、代码重复且难以维护。使用多路输出一次遍历即可完成分流每个输出流各自携带完整的元素和窗口信息后续可以分别继续处理。在 Go SDK 中ParDo输出数量的不同对应不同的 API输出数量API返回值1beam.ParDobeam.PCollection2beam.ParDo2(beam.PCollection, beam.PCollection)3beam.ParDo33 个beam.PCollection4beam.ParDo44 个beam.PCollection5beam.ParDo55 个beam.PCollection6beam.ParDo66 个beam.PCollection7beam.ParDo77 个beam.PCollection二、Kata 实战为大于 100 的数字增加额外输出本课的目标是完成这样一个练习Kata为你的 ParDo 实现额外输出把大于 100 的数字输出到第二个 PCollection。练习所在的目录是 learning/katas/go/core_transforms/additional_outputs/additional_outputs其标准工程结构如下additional_outputs/ ├── cmd/ │ └── main.go # 完整管道入口可见 ├── pkg/ │ └── task/ │ └── task.go # 待补全的 TODO 占位可见 ├── test/ │ └── task_test.go # 测试用例隐藏 ├── task-info.yaml ├── task-remote-info.yaml └── task.md # 本课任务说明从 task-info.yaml 可以看出test/task_test.go对学习者隐藏而cmd/main.go与pkg/task/task.go可见——其中task.go内有一段TODO()占位符偏移 972、长度 164正是需要你补全实现的位置。DoFn 的签名要求使用beam.ParDo2输出两个 PCollection 时DoFn 必须采用如下签名func doFn(element X, emit1 func(Y), emit2 func(Y)) { // element 类型为 X来自输入 PCollection // 调用 emit1 向第一个输出 PCollection 发射元素 // 调用 emit2 向第二个输出 PCollection 发射元素 }关键点在于DoFn 的入参除了元素本身还包含两个或更多emit 函数参数。每个 emit 函数对应一个输出Beam 运行时会为每个输出 PCollection 提供独立的发射器。你只需在函数体内按业务规则决定调用哪个 emitBeam 就会把元素分别路由到对应的输出 PCollection。三、完整解决方案3.1 核心变换pkg/task/task.go本课的标准解法实现如下对应仓库中的 task.gopackage task import ( github.com/apache/beam/sdks/v2/go/pkg/beam ) 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) }分析这段实现ApplyTransform接收beam.Scope变换作用域和输入beam.PCollection返回两个输出 PCollectionbeam.ParDo2的第二个参数是一个匿名 DoFn 函数func(element int, numBelow100, numAbove100 func(int))——第一个参数element int是输入元素后两个参数是输出发射器业务分流逻辑非常直观element 100时调用numBelow100(element)否则调用numAbove100(element)每个元素只会进入两个输出流之一实现了干净的一对多路由。3.2 完整管道cmd/main.gomain.go 给出了可独立运行的完整管道package main import ( beam.apache.org/learning/katas/core_transforms/additional_outputs/additional_outputs/pkg/task context github.com/apache/beam/sdks/v2/go/pkg/beam github.com/apache/beam/sdks/v2/go/pkg/beam/log github.com/apache/beam/sdks/v2/go/pkg/beam/x/beamx github.com/apache/beam/sdks/v2/go/pkg/beam/x/debug ) func main() { ctx : context.Background() p, s : beam.NewPipelineWithRoot() input : beam.Create(s, 10, 50, 120, 20, 200, 0) numBelow100, numAbove100 : task.ApplyTransform(s, input) debug.Printf(s, Number 100: %v, numBelow100) debug.Printf(s, Number 100: %v, numAbove100) err : beamx.Run(ctx, p) if err ! nil { log.Exitf(context.Background(), Failed to execute job: %v, err) } }管道结构拆解创建管道beam.NewPipelineWithRoot()返回管道对象p和根作用域s后续所有变换都挂载在s下构造输入beam.Create(s, 10, 50, 120, 20, 200, 0)创建一个包含 6 个整数的内存 PCollection其中 10、50、20、0 属于 100一类120、200 属于 100一类应用多输出变换task.ApplyTransform(s, input)同时拿到numBelow100和numAbove100两条输出流打印结果debug.Printf是 Beam Go SDK 提供的调试变换位于beam/x/debug包在管道执行时把 PCollection 中的元素格式化打印执行管道beamx.Run(ctx, p)在默认 runner 上运行管道出错时通过log.Exitf记录日志并退出。运行该程序后预期输出类似Number 100: [0 10 20 50] Number 100: [120 200]注意由于debug.Printf处理的是无界意义上的分布式 PCollection元素顺序不保证与输入顺序一致但集合内容应与上述一致。四、测试验证test/task_test.goBeam Go SDK 内置了轻量的测试框架本课使用passertPCollection 断言和ptest测试运行器来验证多输出是否正确。task_test.go 的核心测试如下func TestApplyTransform(t *testing.T) { p, s : beam.NewPipelineWithRoot() tests : []struct { input beam.PCollection wantBelow100 []interface{} wantAbove100 []interface{} }{ { input: beam.Create(s, 10, 50, 120, 20, 200, 0), wantBelow100: []interface{}{0, 10, 20, 50}, wantAbove100: []interface{}{120, 200}, }, } for _, tt : range tests { gotBelow100, gotAbove100 : task.ApplyTransform(s, tt.input) passert.Equals(s, gotBelow100, tt.wantBelow100...) passert.Equals(s, gotAbove100, tt.wantAbove100...) if err : ptest.Run(p); err ! nil { t.Error(err) } } }要点测试通过passert.Equals分别断言两条输出流的内容验证分流结果完全符合预期0,10,20,50进入below100120,200进入above100ptest.Run(p)在本地以 Direct runner 语义执行管道无需外部集群由于测试文件对学习者隐藏写完task.go后运行go test ./...即可自动校验实现是否正确。五、ParDo2 的底层原理5.1 从源码看 ParDo 系列在 Go SDK 的 pardo.go 中ParDo2的实现非常简洁// ParDo2 inserts a ParDo with 2 outputs into the pipeline. func ParDo2(s Scope, dofn any, col PCollection, opts ...Option) (PCollection, PCollection) { ret : MustN(TryParDo(s, dofn, col, opts...)) if len(ret) ! 2 { panic(formatParDoError(dofn, len(ret), 2)) } return ret[0], ret[1] }而单输出的ParDo则写成func ParDo(s Scope, dofn any, col PCollection, opts ...Option) PCollection { ret : MustN(TryParDo(s, dofn, col, opts...)) if len(ret) ! 1 { panic(formatParDoError(dofn, len(ret), 1)) } return ret[0] }从源码结构可以看出几个关键设计统一入口TryParDoParDo1ParDo7都复用了同一个底层函数TryParDo见 pardo.go。它负责校验作用域、解析 side input 与 DoFn、处理输入类型普通类型、KV 类型、CoGBK 类型分别设置不同的主输入方式graph.NumMainInputs最终返回一个[]PCollection切片输出数量自检TryParDo返回的切片长度取决于 DoFn 中声明的 emit 函数个数。MustN对错误直接 panic随后len(ret)必须与所选 API 期望的输出数量一致否则触发formatParDoError友好的错误提示formatParDoError会根据 DoFn 实际输出数量与调用 API 期望数量的差异推荐正确的 API例如DoFn name has 3 outputs, but ParDo2 requires 2 outputs, use ParDo3 instead.这意味着如果你给ParDo2传了一个带 3 个 emit 参数的 DoFn程序会在管道构造阶段直接报出可读的错误并提示改用ParDo3而不是在运行时静默出错。5.2 运行时如何分流从 Beam 的编程模型看多路输出的核心是DoFn 的每个 emit 函数在运行时被绑定到对应的输出 PCollection。Beam runner 会把带 N 个 emit 参数的 DoFn 解释为 N 路输出变换每个输出拥有独立的编码器coder与窗口策略。你在 DoFn 内调用哪个 emit元素就被投递到哪个输出流从而实现一次遍历、多路分流。5.3 与 Fusion 优化的关系pardo.go 的文档还揭示了多路输出与 runner 优化的关系Beam runner 会对 ParDo 执行fusion融合优化——若一个 ParDo 的输出仅被另一个 ParDo 消费两者会融合成单个 ParDo 单遍执行producer-consumer fusion若多个 ParDo 共享同一输入 PCollection则融合成一次遍历sibling fusion。这意味着你可以放心地把多路输出拆成模块化、可组合的 ParDo 风格runner 会自动扁平化为高效执行阶段这也正是多输出变换适合做数据分流的原因之一。六、扩展到更多输出ParDo3 到 ParDo7ParDo2只是多输出家族的一员。同一文件中ParDo3ParDo7的实现模式完全一致只是输出数量不同例如// ParDo3 inserts a ParDo with 3 outputs into the pipeline. func ParDo3(s Scope, dofn any, col PCollection, opts ...Option) (PCollection, PCollection, PCollection) { ret : MustN(TryParDo(s, dofn, col, opts...)) if len(ret) ! 3 { panic(formatParDoError(dofn, len(ret), 3)) } return ret[0], ret[1], ret[2] }对应的 DoFn 签名自然扩展为func doFn(element X, emit1 func(Y), emit2 func(Y), emit3 func(Y)) { // 按业务规则分别调用 emit1 / emit2 / emit3 }如果确实需要超过 7 路输出从源码结构看TryParDo本身支持任意数量的 emit[]PCollection切片但 API 层目前提供到ParDo7为止。绝大多数业务场景如按类别、级别、区间分流用 25 路输出已经足够。七、如何运行本课练习7.1 使用 GoLand EduTools推荐本课程设计为 GoLand 配合 EduTools 插件的交互式 Code Kata。learning/katas/go/README.md 给出了标准设置步骤在 GoLand 中通过 EduTools 插件选择 New Project弹出提示时选择Yes从现有源码创建项目等待索引完成打开Project工具窗口切换到Course视图项目即准备就绪可以直接在task.go的TODO()处编写实现并通过测试反馈验证。7.2 命令行方式不依赖 IDE 时也可以直接在learning/katas/go目录下运行# 运行核心测试验证 task.go 实现 go test ./core_transforms/additional_outputs/... # 运行完整管道 go run ./core_transforms/additional_outputs/additional_outputs/cmd前提是本机已安装与仓库go.mod匹配的 Go 版本并已正确配置GOPATH/GO111MODULE等 Go 环境。仓库根目录的gradlew也可以驱动相关 Gradle 任务如rat许可证检查但日常练习直接使用go命令即可。小结通过本课你掌握了 Apache Beam Go SDK 中 Additional Outputs 的完整用法ParDo系列ParDo、ParDo2ParDo7允许一个变换产生多路 PCollection 输出多路输出的 DoFn 通过多个 emit 函数参数声明输出每个 emit 对应一条输出流以大于 100 的数字单独输出为场景的完整代码包含 task.go 的实现与 main.go 的管道骨架测试通过passert与ptest对每路输出分别断言见 task_test.go底层上pardo.go 中TryParDo统一解析 DoFn 的 emit 数量ParDo2等 API 负责数量自检并提供友好的错误提示runner 的 fusion 优化则让多路输出保持模块化的同时高效执行。下次需要把数据按规则分流时不妨直接选用ParDo2/ParDo3一次遍历即可优雅地完成多路输出。赞分享【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam18/beam点击查看免费下载相关推荐Apache Beam Go Katas使用 ParDo 附加输出Additional Outputs一次生成多条 PCollectionApache Beam Go Katas使用 ParDo 附加输出Additional Outputs一次生成多条 PCollection 导读 本篇围绕大数据批处理流处理数据工程Argo CD CLI 命令补全实战指南argocd completion 命令详解与源码原理Argo CD CLI 命令补全实战指南argocd completion 命令详解与源码原理 导读 argocd completion 是 Argo CD大数据批处理流处理数据工程Apache Beam Go SDK 分支Branching实战一个 PCollection 派生多个输出Apache Beam Go SDK 分支Branching实战一个 PCollection 派生多个输出 本指南基于仓库中的 Beam Katas 练习大数据批处理流处理数据工程上一篇如何轻松掌控游戏窗口SRWE终极分辨率调整指南 下一篇LeRobot训练可视化终极指南3步解决机器人模型黑箱难题创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表