【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam18/beam点击查看免费下载本文以 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: KVstring,int/GW/KVbytes,int[varintz]}] [{10: KVint,string/GW/KVint[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: KVstring,int/GW/KVbytes,int[varintz]} {7: CoGBKstring,int/GW/CoGBKbytes,int[varintz]} {8: KVstring,int/GW/KVbytes,int[varintz]} {9: string/GW/bytes} {10: KVint,string/GW/KVint[varintz],bytes} {11: CoGBKint,string/GW/CoGBKint[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: KVT,int - {6: KVstring,int/GW/KVbytes,int[varintz]}] 7: CoGBK [In(Main): KVstring,int - {6: KVstring,int/GW/KVbytes,int[varintz]}] - [Out: CoGBKstring,int - {7: CoGBKstring,int/GW/CoGBKbytes,int[varintz]}] 8: Combine [In(Main): int - {7: CoGBKstring,int/GW/CoGBKbytes,int[varintz]}] - [Out: KVstring,int - {8: KVstring,int/GW/KVbytes,int[varintz]}] 9: ParDo [In(Main): KVstring,int - {8: KVstring,int/GW/KVbytes,int[varintz]}] - [Out: string - {9: string/GW/bytes}] 10: ParDo [In(Main): T - {9: string/GW/bytes}] - [Out: KVint,T - {10: KVint,string/GW/KVint[varintz],bytes}] 11: CoGBK [In(Main): KVint,string - {10: KVint,string/GW/KVint[varintz],bytes}] - [Out: CoGBKint,string - {11: CoGBKint,string/GW/CoGBKint[varintz],bytes}] 12: ParDo [In(Main): CoGBKint,string - {11: CoGBKint,string/GW/CoGBKint[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 raild: 1 ford: 1 bleeds: 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函数签名直接决定管道形状——两个入参、一个返回值即表示对KVstring,int的 PCollection 操作并输出stringBeam 会在构建期做类型检查。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.txtwordcount引入自定义管道选项、静态 DoFn、复合变换与 Runner 选择debugging_wordcount演示日志与指标调试技巧windowed_wordcount引入窗口概念适合流式场景。此外仓库还提供 cookbookcombine/filter/join/max 等、Kafka、Pub/Sub、Avro、Splittable DoFn、跨语言xlang等大量可运行示例均可作为学习与二次开发起点。在 Dataflow Runner 上运行要在 Google Cloud Dataflow 上运行 wordcount指定 Runner 与云相关参数即可$ go run wordcount.go --runnerdataflow --projectYOUR_GCP_PROJECT --regionYOUR_GCP_REGION --staging_locationYOUR_GCS_LOCATION/staging --worker_harness_container_imageYOUR_SDK_HARNESS_IMAGE_LOCATION --outputYOUR_GCS_LOCATION/output各参数含义参数说明--runnerdataflow选择 Dataflow Runner--runner标志由beamx包注册默认 prism--project你的 GCP 项目 ID--region运行 Dataflow Job 的区域--staging_locationGCS 暂存路径用于上传管道与 SDK 相关构件--worker_harness_container_imageGo 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 Crownd: 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 ModulesBeam 的 Go 代码在仓库sdks目录下维护单一 Go Modulesdks/go.modmodule 名为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、PCollectionpcollection.go、ParDopardo.go、Combinecombine.go、Flatten、Partition、指标metrics.go、窗口windowing.go等 APIpkg/beam/io 提供 textio 等 I/Opkg/beam/transforms 提供 stats 等常用变换。对 SDK 开发者克隆仓库后在模块目录repo/sdks下任意子目录内即可开发与测试Go 工具链照常工作。两点注意修改.proto文件后需要重新生成代码参考pkg/beam/model/PROTOBUF.md修改.tmpl文件后需要把 specialize 工具加入 PATHgo 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 RunnerULRUniversal 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/$GOVERSlatest把指定版本安装到$GOPATH/bin下输出GOCMD供后续使用run_with_go_version.sh默认GOVERSgo1.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示例并理解其管道结构用--runnerdataflow及配套参数把管道部署到 Google Cloud理解单一 Go Modulegithub.com/apache/beam/sdks/v2的工程布局与 GoGradle 构建集成最后以go test ./...和./gradlew :sdks:go:...验证你自己的 SDK 改动。更完整的构建细节可继续阅读 sdks/go/BUILD.md。赞分享【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam18/beam点击查看免费下载相关推荐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上一篇Deepspeed分布式训练实战FAQ_Of_LLM_Interview中的Zero优化下一篇react-native-firebase 单仓构建基准benchmark-prepare.sh 的前后对比与 Nx 本地缓存收益分析创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考