【免费下载链接】beam
Apache Beam is a unified programming model for Batch and Streaming data processing.
本文以 Apache Beam 仓库中的 Go SDK 说明文档 为核心骨架,结合仓库内真实的示例代码、构建配置与测试任务,系统讲解 Apache Beam Go SDK 的定位、示例运行方式、Dataflow 云环境执行、Go Modules 工程结构、Gradle 构建集成以及 SDK 源码开发与测试方法。读完本文,你将能够在本地直接运行 Beam Go 管道,将其部署到 Google Cloud Dataflow,并像 Jenkins/CI 一样对 SDK 进行单元测试与多 Runner 验证。
Go SDK 是什么
Apache Beam 提供统一的批处理与流处理编程模型(Unified Programming Model for Batch and Streaming data processing),而Go SDK 正是该模型在 Go 编程语言中的实现。它基于 Beam 社区最初的设计 RFC(见 sdks/go/README.md)发展而来,让 Go 开发者可以用原生 Go 代码编写 Beam 管道:构建 Pipeline、定义 DoFn、应用 ParDo/Combine 等变换、读写文本与云存储,并在不同的 Runner 上执行——从本地的 Prism/direct Runner 到云端的 Google Cloud Dataflow,再到 Flink、Spark 等分布式 Runner。
仓库中 sdks/go 目录就是 Go SDK 的全部代码,其中 pkg/beam 是核心 SDK 包,examples 提供了从入门到进阶的完整示例,BUILD.md 则专门描述 Go 代码布局与 Gradle 构建集成。
如何运行 Go SDK 示例
前置条件
Beam 示例默认读取/写入 Google Cloud 的源与汇(GCS、Pub/Sub 等)。如果要使用这些 Google Cloud 能力,需要先完成 Google Cloud Dataflow 运行环境的设置(启用相应 API、配置凭据、开通 GCS 存储桶等),并可通过先运行对应的 Java 示例来验证环境是否就绪。仅在本机直接 Runner 上跑通示例则不需要云环境。
示例就是普通 Go 程序
Go SDK 的示例与一般 Go 程序无异,绝大多数可以直接运行,它们通过 Go 标准库flag以命令行参数进行配置。以 wordcount 示例 为例,在仓库的sdks/go目录下直接运行:
$ pwd [...]/sdks/go $ go run examples/wordcount/wordcount.go --output=/tmp/result.txt [{6: KV<string,int>/GW/KV<bytes,int[varintz]>}] [{10: KV<int,string>/GW/KV<int[varintz],bytes>}] 2018/03/21 09:39:03 Pipeline: 2018/03/21 09:39:03 Nodes: {1: []uint8/GW/bytes} {2: string/GW/bytes} {3: string/GW/bytes} {4: string/GW/bytes} {5: string/GW/bytes} {6: KV<string,int>/GW/KV<bytes,int[varintz]>} {7: CoGBK<string,int>/GW/CoGBK<bytes,int[varintz]>} {8: KV<string,int>/GW/KV<bytes,int[varintz]>} {9: string/GW/bytes} {10: KV<int,string>/GW/KV<int[varintz],bytes>} {11: CoGBK<int,string>/GW/CoGBK<int[varintz],bytes>} Edges: 1: Impulse [] -> [Out: []uint8 -> {1: []uint8/GW/bytes}] 2: ParDo [In(Main): []uint8 <- {1: []uint8/GW/bytes}] -> [Out: T -> {2: string/GW/bytes}] 3: ParDo [In(Main): string <- {2: string/GW/bytes}] -> [Out: string -> {3: string/GW/bytes}] 4: ParDo [In(Main): string <- {3: string/GW/bytes}] -> [Out: string -> {4: string/GW/bytes}] 5: ParDo [In(Main): string <- {4: string/GW/bytes}] -> [Out: string -> {5: string/GW/bytes}] 6: ParDo [In(Main): T <- {5: string/GW/bytes}] -> [Out: KV<T,int> -> {6: KV<string,int>/GW/KV<bytes,int[varintz]>}] 7: CoGBK [In(Main): KV<string,int> <- {6: KV<string,int>/GW/KV<bytes,int[varintz]>}] -> [Out: CoGBK<string,int> -> {7: CoGBK<string,int>/GW/CoGBK<bytes,int[varintz]>}] 8: Combine [In(Main): int <- {7: CoGBK<string,int>/GW/CoGBK<bytes,int[varintz]>}] -> [Out: KV<string,int> -> {8: KV<string,int>/GW/KV<bytes,int[varintz]>}] 9: ParDo [In(Main): KV<string,int> <- {8: KV<string,int>/GW/KV<bytes,int[varintz]>}] -> [Out: string -> {9: string/GW/bytes}] 10: ParDo [In(Main): T <- {9: string/GW/bytes}] -> [Out: KV<int,T> -> {10: KV<int,string>/GW/KV<int[varintz],bytes>}] 11: CoGBK [In(Main): KV<int,string> <- {10: KV<int,string>/GW/KV<int[varintz],bytes>}] -> [Out: CoGBK<int,string> -> {11: CoGBK<int,string>/GW/CoGBK<int[varintz],bytes>}] 12: ParDo [In(Main): CoGBK<int,string> <- {11: CoGBK<int,string>/GW/CoGBK<int[varintz],bytes>}] -> [] 2018/03/21 09:39:03 Reading from gs://apache-beam-samples/shakespeare/kinglear.txt 2018/03/21 09:39:04 Writing to /tmp/result.txt注意几点:
- 当前的调试输出相当冗长且内容在未来可能调整,这是正常现象,不影响结果文件。
- 上面是在本地直接 Runner 上执行,因此输出为本地文件:
$ head /tmp/result.txt while: 2 darkling: 1 rail'd: 1 ford: 1 bleed's: 1 hath: 52 Remain: 1 disclaim: 1 sentence: 1 purse: 6理解 wordcount 的管道结构
从 wordcount.go 的源码可以看到 Go SDK 管道的标准编写模式:
自定义管道选项:管道选项就是标准 Go flag,例如
input(默认gs://apache-beam-samples/shakespeare/kinglear.txt,可换成其他文件或 glob)和必填的output;此外还有--small_word_length(默认 9)控制"小词"判定阈值。配置管道就是普通 Go 代码,无特殊框架约束。DoFn 的两种形态:
- 结构体 DoFn(如
extractFn):实现ProcessElement(ctx, line, emit func(string))方法,可通过 emit 函数输出任意数量的元素,适合一对多变换; - 函数式 DoFn(如
formatFn(w string, c int) string):函数签名直接决定管道形状——两个入参、一个返回值即表示对KV<string,int>的 PCollection 操作并输出string,Beam 会在构建期做类型检查。
- 结构体 DoFn(如
DoFn 注册:为了让可移植 Runner 在 Worker 上访问这些函数,必须在
init()中注册。结构体 DoFn 使用register.DoFn3x0context.Context, string, func(string)(3 入参 0 出参),函数 DoFn 使用register.Function2x1(formatFn),发射器用register.Emitter1[string]()注册以加速执行。复合变换:
CountWords是一个典型的复合 PTransform——一个把 ParDo 与stats.Count打包成可复用函数的普通 Go 函数,通过s.Scope("CountWords")获得作用域命名,便于监控。运行入口:
beam.Init()必须在启动时调用(在分布式 Runner 上用于接管控制),然后构建 Pipeline、挂上textio.Read/ ParDo / Count /textio.Write等变换,最后用beamx.Run执行——它通过--runner标志选择 Runner,默认是 prism。
一组由浅入深的 wordcount 系列
仓库 examples 中有一组刻意设计的 wordcount 系列示例,帮助循序渐进理解 Beam 概念:
- minimal_wordcount:无参数、无错误处理,聚焦管道构建本身,直接以
prism.Execute(context.Background(), p)在 Prism Runner 上执行,结果写入当前目录wordcounts.txt; - wordcount:引入自定义管道选项、静态 DoFn、复合变换与 Runner 选择;
- debugging_wordcount:演示日志与指标调试技巧;
- windowed_wordcount:引入窗口概念,适合流式场景。
此外仓库还提供 cookbook(combine/filter/join/max 等)、Kafka、Pub/Sub、Avro、Splittable DoFn、跨语言(xlang)等大量可运行示例,均可作为学习与二次开发起点。
在 Dataflow Runner 上运行
要在 Google Cloud Dataflow 上运行 wordcount,指定 Runner 与云相关参数即可:
$ go run wordcount.go --runner=dataflow --project=<YOUR_GCP_PROJECT> --region=<YOUR_GCP_REGION> --staging_location=<YOUR_GCS_LOCATION>/staging --worker_harness_container_image=<YOUR_SDK_HARNESS_IMAGE_LOCATION> --output=<YOUR_GCS_LOCATION>/output各参数含义:
| 参数 | 说明 |
|---|---|
--runner=dataflow | 选择 Dataflow Runner(--runner标志由beamx包注册,默认 prism) |
--project | 你的 GCP 项目 ID |
--region | 运行 Dataflow Job 的区域 |
--staging_location | GCS 暂存路径,用于上传管道与 SDK 相关构件 |
--worker_harness_container_image | Go SDK Harness 容器镜像位置 |
--output | 结果输出的 GCS 路径 |
此时输出是 GCS 文件,可用gsutil查看:
$ gsutil cat <YOUR_GCS_LOCATION>/output* | head Blanket: 1 blot: 1 Kneeling: 3 cautions: 1 appears: 4 Deserved: 1 nettles: 1 OSWALD: 53 sport: 3 Crown'd: 1关于 Go SDK Harness 容器镜像的构建与推送方法,可参考 Beam 容器构建文档(运行时环境一节)。
Runner 注册机制背后的实现
为什么--runner能同时识别 Dataflow、Flink、Spark、Samza、Prism 等多个 Runner?答案在 beamx 包:它通过空导入_ "github.com/apache/beam/sdks/v2/go/pkg/beam/runners/..."触发各 Runner 包注册的副作用,同时导入反射优化运行时与 gcs/local 文件系统,默认 Runner 为 prism。仓库 runners 目录 下即可看到 dataflow、direct、dot、flink、prism、samza、spark、universal 等 Runner 实现,管道作者只需import "github.com/apache/beam/sdks/v2/go/pkg/beam/x/beamx"即可获得全部 Runner 支持。
Go 工程结构与 Go Modules
Beam 的 Go 代码在仓库sdks目录下维护单一 Go Module(sdks/go.mod,module 名为github.com/apache/beam/sdks/v2,声明 go 1.20),这样既覆盖用户管道开发所需的全部 Go 代码,也覆盖执行层代码(包括 Java/Python SDK 目录下的容器 bootloader 代码)。之所以不放在仓库根目录,是因为会与既有 vendor 目录产生冲突(见 sdks/go.mod 头部注释)。
对管道作者
在你的go.mod中声明对github.com/apache/beam/sdks/v2的依赖即可使用 Beam,例如直接拉取核心包:
go get github.com/apache/beam/sdks/v2/go/pkg/beam核心包 pkg/beam 提供 Pipeline 构建(pipeline.go)、PCollection(pcollection.go)、ParDo(pardo.go)、Combine(combine.go)、Flatten、Partition、指标(metrics.go)、窗口(windowing.go)等 API;pkg/beam/io 提供 textio 等 I/O;pkg/beam/transforms 提供 stats 等常用变换。
对 SDK 开发者
克隆仓库后,在模块目录(<repo>/sdks下任意子目录)内即可开发与测试,Go 工具链照常工作。两点注意:
- 修改
.proto文件后需要重新生成代码,参考pkg/beam/model/PROTOBUF.md; - 修改
.tmpl文件后需要把 specialize 工具加入 PATH:
go get github.com/apache/beam/sdks/v2/go/cmd/specialize export PATH=$PATH:$GOROOT/bin:$GOPATH/bin与 Gradle 的构建集成
GoGradle 插件
Beam 通过名为 GoGradle 的 Gradle 插件把 Go 代码纳入整体构建,但禁用 GoGradle 的 vendoring,改为使用 Go Modules 管理依赖。GoGradle 负责在 Gradle/Jenkins 中调用 go 工具链,与 SDK 贡献者和用户使用同一套依赖。对于少量 Go 二进制(如容器 bootloader),同一份代码既可用 Gradle 构建,也可用标准 Go 工具构建。
容器镜像构建还带来一个特殊点:镜像通常面向 linux/amd64,而开发机可能不是该架构,因此需要为容器镜像交叉编译 Go 二进制(一般放在target/linux_amd64)。
验证 SDK 的测试任务
在 beam 根目录下,可用与 Jenkins 完全一致的方式验证改动:
./gradlew :sdks:go:goTest—— 执行 SDK 单元测试;./gradlew :sdks:go:test:prismValidatesRunner—— 以独立二进制(含容器)方式验证 SDK 在 Go Prism Runner 上的行为;./gradlew :sdks:go:test:ulrValidatesRunner—— 验证 SDK 在 Portable Python Runner(ULR,Universal Local Runner)上的行为;./gradlew :sdks:go:test:flinkValidatesRunner—— 验证 SDK 在 Flink Runner 上的行为。
同时,在<beam root>/sdks/go目录直接执行go test ./...可运行 SDK 全部单元测试,这正是 README 推荐的日常开发验证方式。
Go 版本管控
仓库提供两个脚本固定 Beam 基础设施使用的 Go 版本:
- prepare_go_version.sh:通过
--version指定完全限定版本(如go1.21.0),利用 Go 1.16+ 的任意版本下载能力,借助go install golang.org/dl/$GOVERS@latest把指定版本安装到$GOPATH/bin下,输出GOCMD供后续使用; - run_with_go_version.sh:默认
GOVERS=go1.21.0,支持--version与--gocmd两个可选标志,在准备完成后以文件锁保证并发安全地执行指定版本的 go 命令(flock排他锁做下载、共享锁做普通执行)。
这两个脚本配合使用,可实现与 Jenkins 一致的、可复现的 hermetic 构建环境。
参与 Go SDK 开发
从零开始的 Go 基础
如果你刚接触 Go,可先通过 Go Tour 交互式学习语言基础(无需安装),再参考 Go 工具链实战工作坊了解推荐的开发工具链,然后回到仓库实操。
开发流程
Go SDK 使用 Go Modules 管理依赖,开发流程就是克隆仓库 → 修改代码 → 运行测试三步,与普通 Go 项目无异:
- 在
sdks/go目录执行go test ./...跑单元测试; - 按上述 Gradle 任务做多 Runner 验证;
- 按 Beam 贡献指南创建分支、提交 Pull Request。
提交前请确保注册了新用到的 DoFn/函数(否则可移植 Runner 无法在 Worker 上调用),并为新逻辑补充测试——仓库 examples 下已有多处*_test.go(如 large_wordcount_test.go、snippets/04transforms_test.go)可以作为参考模板。
问题反馈
任何 bug 或功能需求,请在 issue 中使用sdk-go组件标签进行反馈,以便维护者准确分类与跟进。
小结
Apache Beam Go SDK 让你用纯 Go 编写统一批/流管道,并在从本地 Prism 到云端 Dataflow、再到 Flink/Spark 等 Runner 之间无缝切换。掌握本文内容后,你可以:在本地直接运行go run示例并理解其管道结构;用--runner=dataflow及配套参数把管道部署到 Google Cloud;理解单一 Go Module(github.com/apache/beam/sdks/v2)的工程布局与 GoGradle 构建集成;最后以go test ./...和./gradlew :sdks:go:...验证你自己的 SDK 改动。更完整的构建细节可继续阅读 sdks/go/BUILD.md。
【免费下载链接】beam
Apache Beam is a unified programming model for Batch and Streaming data processing.
相关推荐
Warp 补全引擎的 Basic Parser 架构:从 Lex 到类型驱动 Full Parse 的递归下降解析器全解析
Warp 补全引擎的 Basic Parser 架构:从 Lex 到类型驱动 Full Parse 的递归下降解析器全解析 导读 本文深入剖析 Warp 开源仓
批处理流处理大数据Apache Beam Go SDK 实战指南:示例运行、Dataflow 部署与本地构建测试
Apache Beam Go SDK 实战指南:示例运行、Dataflow 部署与本地构建测试 Apache Beam Go SDK 是 Apache Beam
大数据批处理流处理数据工程Apache Beam TypeScript SDK 开发指南:从源码构建、运行 Pipeline 到可移植运行器的实现原理
Apache Beam TypeScript SDK 开发指南:从源码构建、运行 Pipeline 到可移植运行器的实现原理 导读 本文面向希望以 JavaSc
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考