【免费下载链接】beam
Apache Beam is a unified programming model for Batch and Streaming data processing.
导读
本文围绕 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_onetomany | ParDo 一对多 | 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.go:
package 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 管道。
【免费下载链接】beam
Apache Beam is a unified programming model for Batch and Streaming data processing.
相关推荐
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 是最核心的
大数据批处理流处理数据工程
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考