尧图网络科技YAOTU DIGITAL 获取报价
获取报价
首页 / 资讯中心 / 文章详情

Apache Beam Go 实战:用 ParDo 实现 One-to-Many 一对多映射

发布时间:2026/9/29 7:29:22

资讯中心
01
ARTICLE

Apache Beam Go 实战:用 ParDo 实现 One-to-Many 一对多映射

Apache Beam Go 实战:用 ParDo 实现 One-to-Many 一对多映射
【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam18/beam点击查看免费下载导读本文围绕 Apache Beam Go SDK 的 ParDo 一对多One-to-Many映射模式展开讲解如何在单个输入元素的基础上产生零个、一个或多个输出元素。文中将以「把句子按空格拆分为单词」这一经典 kata 为例完整给出可运行的 Go 代码、测试用例与底层实现原理帮助读者掌握 Go 语言下 ParDo 与 DoFn 的编写范式并理解其与一对一映射的本质区别。ParDo 是什么从 Map 到一对多在 Apache Beam 中ParDo 是用于通用并行处理的核心 PTransform其处理范式与 Map/Shuffle/Reduce 算法中的 Map 阶段类似它逐个考察输入 PCollection 中的每个元素调用用户自定义的处理函数即 DoFn然后向输出 PCollection 发射零个、一个或多个元素。在 Beam 学习训练营Katas中learning/katas/go/core_transforms/map/目录下的课程按难度递进编排见 lesson-info.yaml课程主题映射关系pardo基础 ParDo一对一1 个输入 → 1 个输出pardo_onetomanyParDo 一对多1 个输入 → 多个输出pardo_struct结构体 DoFn使用 struct 形式编写 DoFn上一课pardo中DoFn 是一个纯函数输入一个元素、返回一个元素func multiplyBy10Fn(element int) int { return element * 10 } func ApplyTransform(s beam.Scope, input beam.PCollection) beam.PCollection { return beam.ParDo(s, multiplyBy10Fn, input) }本课pardo_onetomany要解决的是相反的问题一个输入元素如何变成多个输出元素。最直观的场景就是把一个句子按空格拆分成多个单词。实战 Kata把句子拆成单词本课的练习文档见 task.md其练习目标如下请编写一个 ParDo将每个输入句子按空格 切分成单词。骨架与占位符与所有 Beam Katas 一样本课通过 task-info.yaml 定义了练习结构test/task_test.go对学员隐藏visible: false而pkg/task/task.go与cmd/main.go可见其中task.go的两个TODO()占位符就是学员需要补全的位置。入口程序 cmd/main.go 已经搭好了整条流水线func main() { p, s : beam.NewPipelineWithRoot() input : beam.Create(s, Hello Beam, It is awesome) output : task.ApplyTransform(s, input) debug.Print(s, output) err : beamx.Run(context.Background(), p) if err ! nil { log.Exitf(context.Background(), Failed to execute job: %v, err) } }它做了三件事用beam.Create创建包含两个字符串Hello Beam与It is awesome的输入 PCollection调用task.ApplyTransform施加自定义变换用beamx.Run在 Direct Runner 上执行并用debug.Print输出结果。参考答案DoFn 配合 emit 回调一对多映射的关键在于DoFn 不再返回单个值而是通过一个emit回调函数逐条发射结果。完整实现见 pkg/task/task.gopackage task import ( github.com/apache/beam/sdks/v2/go/pkg/beam strings ) func ApplyTransform(s beam.Scope, input beam.PCollection) beam.PCollection { return beam.ParDo(s, tokenizeFn, input) } func tokenizeFn(input string, emit func(out string)) { tokens : strings.Split(input, ) for _, k : range tokens { emit(k) } }逐段拆解tokenizeFn(input string, emit func(out string))第一个参数是输入元素类型第二个参数emit是输出回调。DoFn 内每调用一次emit(k)就会向输出 PCollection 发射一个元素。strings.Split(input, )按单个空格切分句子得到单词切片。循环调用emit(k)把每个单词逐一发射出去实现「一个句子 → N 个单词」的一对多映射。运行该程序输入Hello Beam和It is awesome会被展开为五个单词Hello、Beam、It、is、awesome。测试验证隐藏的测试文件 test/task_test.go 给出了标准断言func TestTask(t *testing.T) { p, s : beam.NewPipelineWithRoot() tests : []struct { input beam.PCollection want []interface{} }{ { input: beam.Create(s, Hello Beam. It is awesome.), want: []interface{}{Hello, Beam., It, is, awesome.}, }, } 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) } } }这里用passert.Equals对变换结果与期望序列逐元素比对用ptest.Run在内存中执行整个 pipeline。注意want保留了句点Beam.、awesome.是原单词的一部分因为按 单个空格切分不会剥离标点。这提示了一个进阶问题——真实场景往往还需要过滤空串或去除标点可在 DoFn 中自行扩展。深入原理Go SDK 中 ParDo 是如何工作的一对多模式在 Beam Go SDK 中是 ParDo 的内建能力。在 sdks/go/pkg/beam/pardo.go 中ParDo的文档明确说明ParDo 是 Apache Beam 中核心的逐元素 PTransform对输入 PCollection 的每个元素调用用户指定函数产生零个或多个输出元素全部收集到输出 PCollection 中。DoFn 的两种形态从源码 sdks/go/pkg/beam/pardo.go#L153-L162 可以看到DoFn 有两种写法单个函数如本课的tokenizeFn(input string, emit func(out string))。Go SDK 通过反射识别函数签名若函数第二个参数是func(out T)形式的回调则该 DoFn 支持一对多flatMap 语义发射若函数仅返回一个值则是一对一映射。结构体struct实现ProcessElement等方法并可选实现Setup、StartBundle、FinishBundle、Teardown生命周期方法见下一课pardo_struct。注册与序列化约束源码中同时强调了两条关键约束DoFn 必须是包级具名函数不能是匿名函数或闭包否则在分布式 worker 上执行时会失败用作 DoFn 的函数与类型必须通过beam的register包注册以便在分布式执行时序列化分发。这意味着在本课这类只跑 Direct Runner 的本地练习中可以不注册但部署到 Dataflow、Flink 等分布式 Runner 时register.Function1x1/register.DoFn等注册步骤是必不可少的。一对多时的内部行为从 sdks/go/pkg/beam/pardo.go#L428-L434 可以看到ParDo是TryParDo的便捷封装它要求 DoFn 恰好产生 1 个输出 PCollection否则会 panic。而ParDoN多输出、ParDo2/ParDo3固定多输出等变体则用于需要发射到多个 PCollection 的场景。无论哪种变体单个输入元素「零个或多个输出」的能力都由 DoFn 签名决定这正是 ParDo 比单纯Map更灵活的根源——它天然覆盖了 filter零输出、map单输出、flatMap多输出三种语义。小结一对多映射的适用场景与学习路径一句话总结当 DoFn 的签名包含emit func(T)回调时ParDo 就从「一对一」升级为「一对多」。这一模式在 Go 语言中对应 flatMap/explode 语义广泛用于文本分词本课场景句子 → 单词、日志 → 字段数据展开JSON 数组 → 多条记录、嵌套结构 → 扁平行过滤与转换混合只发射满足条件的元素零输出即等价于过滤。完成本课练习后建议按 lesson-info.yaml 继续pardo_struct课程学习用结构体 DoFn 携带构造期配置、管理有状态资源从而写出更贴近生产环境的 Beam Go 管道。赞分享【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam18/beam点击查看免费下载相关推荐Apache Beam Java Kata 实战用 FlatMapElements 实现一对多one-to-many映射Apache Beam Java Kata 实战用 FlatMapElements 实现一对多one to many映射 FlatMapElements大数据批处理流处理数据工程Apache Beam Go SDK 实战用 ParDo 实现 One-to-Many 一对多变换句子分词 Kata 详解Apache Beam Go SDK 实战用 ParDo 实现 One to Many 一对多变换句子分词 Kata 详解 Apache Beam 的 P大数据批处理流处理数据工程Apache Beam Kotlin Katas 实战用 ParDo 实现 OneToMany 一对多映射Apache Beam Kotlin Katas 实战用 ParDo 实现 OneToMany 一对多映射 Apache Beam 的 ParDo 是最核心的大数据批处理流处理数据工程上一篇抖音批量下载五步搞定抖音视频下载工具 douyin-downloader 实战指南下一篇Wand(WeMod) 2小时限制破解完整教程Wand-Enhancer本地免费补丁手机远程面板实测创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
02
RELATED NEWS

相关资讯

更多网站建设与数字化升级内容

03
WHY YAOTU

想打造同款高转化官网?

懂行业、懂生意,从建站到增长一站式陪跑

◈

场景化定制

不做模板站,围绕你的业务场景量身设计,小众不撞款。

◐

营销型架构

以转化目标组织内容与路径,让官网真正带来询盘。

▲

全周期服务

设计、开发、运营、运维一体,上线只是开始。

免费获取你的建站方案

留下需求,专属顾问 24 小时内为你输出方案建议。