新闻详情

新闻详情

首页 / 资讯中心 / 详情

Apache Beam Go SDK ParDo 入门实战:用 Go 编写并行元素级转换(Multiply by 10 Kata 全解析)

发布时间:2026/9/26 10:37:08来源:尧图网络
Apache Beam Go SDK ParDo 入门实战:用 Go 编写并行元素级转换(Multiply by 10 Kata 全解析)
大数据批处理流处理数据工程【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam4/beam点击查看免费下载本篇技术指南围绕 Apache Beam Go SDK 中最重要的核心转换PTransform——ParDo展开。ParDo 是 Beam 用于通用并行处理的基础原语其处理范式与 Map/Shuffle/Reduce 风格算法中的 Map 阶段类似它会逐个处理输入 PCollection 中的每个元素调用你编写的用户处理函数并零个、一个或多个地输出到目标 PCollection。本文将以仓库中 Go SDK Katas 课程里 map/pardo 这道经典练习题将输入元素乘以 10为主线讲解 ParDo 的概念模型、DoFn 的编写方式、完整可运行代码、测试验证方法以及从 Go 源码层面对 ParDo 执行机制与变体的深入剖析读完后你将能独立用 Go 编写并测试自己的 ParDo 转换。课程背景与任务定位learning/katas/go是 Apache Beam 仓库中的 Go SDK Code Katas代码练习课程采用 GoLand EduTools 插件以交互式练习的形式组织包含 Introduction、Core Transforms、Common Transforms、Windowing、IO 等模块见 learning/katas/go/README.md。本任务位于Core Transforms → Map → ParDo是核心转换课程的第一课与同目录下的pardo_onetomany一对多输出、pardo_struct使用结构体 DoFn构成由浅入深的 ParDo 学习序列其课程内容编排可见 learning/katas/go/core_transforms/map/lesson-info.yaml。本任务的题目Kata编写一个简单的 ParDo将输入元素乘以 10。这是一个典型的 1 对 1one-to-one映射场景每个输入元素恰好产生一个输出元素与Map语义完全一致。虽然 ParDo 的能力远不止映射它可以过滤、聚合、一对多输出但这一课正是理解其最小可用形态的最佳起点。ParDo 的概念模型Beam 的 Map 阶段在动手写代码之前先建立正确的概念模型。ParDo 在 Beam 中承担的角色相当于 MapReduce 风格算法中的 Mapper逐元素处理ParDo 将输入 PCollection 中的每一个元素视为独立处理单元用户代码介入处理逻辑你的业务函数完全由用户编写Beam 框架负责调度与分发零到多输出一个输入元素可以产生 0 个、1 个或多个输出元素全部汇入输出 PCollection分布式并行元素彼此独立处理可以在分布式集群上并行执行。在 Go SDK 中beam.ParDo的官方注释对这一语义有精确描述ParDo is the core element-wise PTransform in Apache Beam, invoking a user-specified function on each of the elements of the input PCollection to produce zero or more output elements, all of which are collected into the output PCollectionParDo 是 Apache Beam 中的核心逐元素 PTransform它对输入 PCollection 的每个元素调用用户指定的函数产生零个或多个输出元素全部收集到输出 PCollection 中并明确说明其处理风格与 MapReduce 中的 Mapper 或 Reducer 类相似见 pardo.go。概念上ParDo 执行时输入元素会被划分为若干 bundle批次分发到分布式工作节点或本地 runner 实例上并行处理每个元素携带的时间戳与所在窗口会原样传递给输出元素。编写第一个 ParDoMultiply by 10 的两种写法方式一普通函数作为 DoFn本任务的标准答案Go SDK 中最简单的 DoFn 就是一个普通函数。Kata 的标准解法位于 learning/katas/go/core_transforms/map/pardo/pkg/task/task.gopackage task import github.com/apache/beam/sdks/v2/go/pkg/beam func ApplyTransform(s beam.Scope, input beam.PCollection) beam.PCollection { return beam.ParDo(s, multiplyBy10Fn, input) } func multiplyBy10Fn(element int) int { return element * 10 }关键点拆解multiplyBy10Fn(element int) int即 DoFn入参是输入元素int返回值是输出元素int。Beam Go SDK 通过反射机制识别这种 1 进 1 出 的函数签名并自动推断其类型与 coder编码器beam.ParDo(s, dofn, input)将 DoFn 应用到输入 PCollection 上返回新的输出 PCollection。第一个参数s beam.Scope是命名作用域用于组织 pipeline 图中的变换层级输出 PCollection 的元素类型由 DoFn 返回值决定这里仍为int数值为输入值的 10 倍。方式二带 emit 回调函数的写法当 DoFn 需要产生多个输出或按条件选择性输出时可在函数签名中增加一个emit回调参数如本课程下一课 learning/katas/go/core_transforms/map/pardo_onetomany/pkg/task/task.go 所示func tokenizeFn(input string, emit func(out string)) { tokens : strings.Split(input, ) for _, k : range tokens { emit(k) } }Go SDK 的 ParDo 支持 0N 个输出 PCollection对应一组变体函数ParDo0、ParDo1 个输出、ParDo2ParDo7多输出以及任意输出数量的ParDoN它们的统一入口是TryParDo定义见 pardo.go。完整可运行示例从 Pipeline 到打印输出要真正运行这个 Kata需要一个完整的 pipeline。仓库中的 learning/katas/go/core_transforms/map/pardo/cmd/main.go 给出了可运行的完整程序package main import ( beam.apache.org/learning/katas/core_transforms/map/pardo/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, 1, 2, 3, 4, 5) output : task.ApplyTransform(s, input) debug.Print(s, output) err : beamx.Run(ctx, p) if err ! nil { log.Exitf(context.Background(), Failed to execute job: %v, err) } }逐行解读这个最小 pipelinebeam.NewPipelineWithRoot()创建新的 Pipeline 并返回根作用域s见 util.go。Go SDK 采用先构图、后执行的模式所有beam.*调用只是把变换节点加入有向无环图DAG真正执行发生在最后一步beam.Create(s, 1, 2, 3, 4, 5)从一组内存值创建输入 PCollection见 create.gotask.ApplyTransform(s, input)应用我们的 ParDo将 1、2、3、4、5 分别乘以 10得到 10、20、30、40、50debug.Print(s, output)是 Go SDK 提供的调试变换把 PCollection 内容打印到日志见 print.gobeamx.Run(ctx, p)在默认 runner 上实际执行整个 pipeline来自x/beamx包。在 Katas 练习环境中默认使用 Direct Runner 在本地执行该 runner 的实现位于 runners/direct-javaJava 实现以及 Go 端对应的本地执行逻辑中若未指定 runnerbeamx.Run会回退到本地直接执行。运行该程序后日志中会输出[10 20 30 40 50]。DoFn 的生命周期与类型要求Go 源码视角从 pardo.go 的文档注释可以提炼出 Go SDK DoFn 的完整规则函数式 DoFn 的限制禁止匿名函数与闭包DoFn 必须是具名函数因为其名称会被用作分布式 worker 上的标识匿名/闭包函数没有稳定名称会在执行期失败必须注册用于 DoFn 的函数和类型必须通过beam的register包注册例如register.Function1x1(fn)这样它们才能被序列化并分发到分布式 worker 上执行。在单机 Direct Runner 下即使不显式注册通常也能运行但生产环境如 Dataflow、Flink runner是必需的。结构体式 DoFn 的生命周期DoFn 也可以是结构体通过定义特定方法参与完整生命周期方法调用时机典型用途Setup每个 worker 创建 DoFn 实例后调用一次初始化非序列化资源如建立数据库连接StartBundle每个 bundle 处理开始前调用初始化本次 batch 处理所需的临时状态ProcessElement对 bundle 中每个输入元素调用核心处理逻辑产生零到多个输出FinishBundle每个 bundle 处理结束后调用冲刷缓冲、聚合本次 batch 结果TeardownDoFn 实例被废弃或异常终止时调用释放资源执行流程为worker 从 JSON 反序列化出全新 DoFn 实例 → 调用Setup→ 对每个 bundle 依次调用StartBundle→ 对 bundle 内每个元素调用ProcessElement→ 调用FinishBundle若任一环节返回错误会触发Teardown。runner 可能复用 DoFn 实例处理多个 bundle但异常终止的实例绝不会被复用。这一机制正是本课程第三课pardo_struct的主题。输出语义所有 DoFn 实例产生的输出元素共同构成输出 PCollection输出元素继承输入元素的时间戳与窗口这是后续 Windowing 课程learning/katas/go/windowing的基础输出 PCollection 之间类型不必相同ParDo2ParDo7正是为此设计。用测试验证 Kata 答案Katas 课程为每个任务提供了隐藏的测试文件用于自动判定答案是否正确。本任务的测试位于 learning/katas/go/core_transforms/map/pardo/test/task_test.gofunc TestApplyTransform(t *testing.T) { p, s : beam.NewPipelineWithRoot() tests : []struct { input beam.PCollection want []interface{} }{ { input: beam.Create(s, 1, 2, 3, 4, 5), want: []interface{}{10, 20, 30, 40, 50}, }, } for _, tt : range tests { got : task.ApplyTransform(s, tt.input) passert.Equals(s, got, tt.want...) if err : ptest.Run(p); err ! nil { t.Error(err) } } }测试模式清晰且可复用通过beam.NewPipelineWithRoot()创建测试 pipeline用beam.Create构造输入调用被测的task.ApplyTransform得到输出用passert.Equals断言输出 PCollection 等于期望值{10, 20, 30, 40, 50}用ptest.Run(p)在 Direct Runner 上执行 pipeline见 ptest.go。若你的实现正确测试通过若multiplyBy10Fn写错比如乘了别的数passert.Equals会给出期望与实际结果的差异。这也是你练习 ParDo 时最快的反馈闭环。ParDo 的进阶形态从一进一出到多输出、带状态掌握基础之后ParDo 的变体可以覆盖几乎所有的逐元素处理需求过滤1→0 或 1→1DoFn 不调用emit即丢弃该元素实现 filter 语义一对多1→N如上面的tokenizeFn示例对一句话按空格切分成多个词依次emit这正是pardo_onetomany一课的内容多输出 PCollection1→N 个 PCollection通过ParDo2ParDo7或ParDoN将元素按条件分流到不同类型/不同语义的输出典型场景如正常数据与异常数据分离旁路输入Side InputParDo可接收额外 PCollection 作为旁路输入beam.SideInput选项实现类似广播 join 的语义见 pardo.go有状态处理State与定时器Timer结构体 DoFn 结合beam.State/beam.Timer参数可实现跨元素的累积与定时触发这是 Go SDK 更高级的用法。这些能力共同说明ParDo 不只是 Map它是 Beam 统一批流模型unified programming model for Batch and Streaming data processing即本项目定位中承载用户业务逻辑的核心载体。小结通过map/pardo这个 Kata你已经掌握了ParDo 的概念模型逐元素、并行、零到多输出的通用处理范式Go SDK 中最简 DoFn 写法具名普通函数一进一出一个完整可运行的 Beam Go pipelineNewPipelineWithRoot→Create→ParDo→debug.Print→beamx.Run结构体 DoFn 的生命周期方法Setup/StartBundle/ProcessElement/FinishBundle/Teardown与函数式 DoFn 的注册、具名约束用passertptest编写 ParDo 单元测试的方法。下一步你可以继续完成同课程下的pardo_onetomany一对多与pardo_struct结构体 DoFn任务并在 learning/katas/go 课程树的core_transforms/map目录下逐个攻克其余转换更深层的 ParDo 实现细节可查阅 pardo.go 的完整注释与实现。赞分享大数据批处理流处理数据工程【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam4/beam点击查看免费下载相关推荐zerotier-cli join 后返回 access denied 时如何确认设备已被控制器授权zerotier cli join 后返回 access denied 时如何确认设备已被控制器授权 在 Unix 系统Linux/BSD/OSX上用 z大数据批处理流处理数据工程theHarvester 提示缺少 API 密钥时怎么处理theHarvester 提示缺少 API 密钥时怎么处理 在 theHarvester 中运行某个发现源discovery source时终端会输出类似大数据批处理流处理数据工程如何用 litgpt serve 开启 OpenAI 兼容端点并让现有 OpenAI SDK 客户端接入本地模型如何用 litgpt serve 开启 OpenAI 兼容端点并让现有 OpenAI SDK 客户端接入本地模型 如果你的应用已经在用 OpenAI API比大数据批处理流处理数据工程上一篇OpenVINO6 框架模型转换到 3 类硬件设备部署下一篇Telegraf enum 处理器插件字段与标签枚举值映射配置与源码实战指南创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
网站建设高端定制企业官网
RELATED

相关资讯

更多精彩内容,欢迎继续阅读

较早相关资讯

最新相关资讯

OA+CRM源码部署与二次开发全流程解析:从架构到避坑 2026/9/26 11:32:18

OA+CRM源码部署与二次开发全流程解析:从架构到避坑

简介:一套面向企业级应用开发者的OA办公系统完整源码包,在组织流程自动化、文档管理、任务协作基础上,额外集成CRM客户管理系统与内部即时聊天工具,并针对手机端做了自适应适配,适合需要学习或二次开发企业协同平台的P…

阅读更多 →
Transformer原理与PyTorch手写实战:从自注意力到LoRA微调 2026/9/26 11:32:18

Transformer原理与PyTorch手写实战:从自注意力到LoRA微调

很多朋友刷到过类似的视频:封面写着“Transformer 从入门到天花板”“保姆级精讲”,点进去弹幕齐刷刷“学会了”,可一关屏幕,连 Positional Encoding 的代码都写不出来。原因不是你不聪明,而是视频节奏太快、信息密度太…

阅读更多 →
Transformer 从原理到实战:手写实现与 LoRA 高效微调 2026/9/26 11:32:18

Transformer 从原理到实战:手写实现与 LoRA 高效微调

直接在浏览器里刷到 Transformer 相关的视频或文章,第一反应往往是“这是一个深度学习基础模型,很重要”,但真到自己动手跑代码时,就会发现网上资料要么只讲论文,要么只贴代码,很少有把“原理拆解→手写实现…

阅读更多 →
OA+CRM+聊天工具源码解析:从部署到移动端适配全攻略 2026/9/26 11:32:18

OA+CRM+聊天工具源码解析:从部署到移动端适配全攻略

简介:这是一套面向企业级应用开发者的OA办公系统完整源码,整合CRM客户管理与内部聊天工具,并针对手机端做了自适应优化,适合需学习企业信息化系统搭建、二次开发或用于毕业设计的开发者。资源包为zip压缩格式,共3272个…

阅读更多 →
podofo 0.9.5 预编译库:VS2013 x86 工程集成 PDF 解析与生成方案 2026/9/26 11:32:17

podofo 0.9.5 预编译库:VS2013 x86 工程集成 PDF 解析与生成方案

简介:面向 Visual Studio 2013 和 x86 平台的 Podofo 0.9.5 预编译库,特别适合在 Windows 下从事 C PDF 开发的工程师。Podofo 是开源且稳定的 PDF 处理库,提供文档读取、解析、修改与生成能力,支持操作页面、字体、加密信息、书签…

阅读更多 →
GEO卫星星点轨迹仿真:从轨道根数到8字曲线的完整计算流程 2026/9/26 11:32:11

GEO卫星星点轨迹仿真:从轨道根数到8字曲线的完整计算流程

简介:面向卫星通信、轨道力学及遥感方向的工程技术人员与高校学生,这份GEO卫星轨迹模拟MATLAB实现包可用于掌握同步轨道卫星相对地面静止的轨迹特点,并辅助开展轨道可视化与通信链路设计。资源共包含4个文件,其中3个.m脚本分别负责…

阅读更多 →

今日资讯

本周资讯

本月资讯

看完文章仍有疑问?

联系尧图顾问,获取一对一建站咨询

立即免费咨询 📞 400-888-8888
📞 ✉