- 大数据
- 批处理
- 流处理
- 数据工程
【免费下载链接】beam
Apache Beam is a unified programming model for Batch and Streaming data processing.
导读
Apache Beam Playground 是 Apache Beam 提供的浏览器端代码运行环境,让用户无需任何本地配置即可直接运行 Beam 管道示例。本文围绕仓库中 playground/TASKS.md 这份开发任务清单展开,系统讲解 Playground 开发与维护过程中最核心的实操内容:预提交检查与 Gradle 任务体系、基于 Docker Compose 的本地全栈环境搭建、Cloud Datastore 中代码片段的清理策略、Beam SDK / Go / Python / Java / Flutter 各语言依赖的升级流程,以及如何将仓库中的示例、Katas 与测试批量部署到 Playground 后端。读完本文,你将掌握从零启动本地 Playground、维护示例数据、同步 SDK 版本到完成一次完整 CD 部署的完整链路。
一、Playground 的 Gradle 任务体系
Playground 的所有构建与运维操作都以 Gradle 任务的形式暴露,统一通过仓库根目录的./gradlew执行。开始之前建议先执行预提交检查,确保改动不会破坏现有构建。
1. 执行整体预提交检查
在仓库根目录运行以下命令,可以对 Playground 全模块执行一次完整的预提交检查:
cd beam ./gradlew playgroundPrecommit该任务是 CI 中针对 Playground 模块的汇总入口,包含编译、单元测试、lint 等检查项,是提交代码前推荐的验证步骤。
2. 查看 Playground 的全部可用任务
cd beam ./gradlew playground:tasks该命令会列出playground工程(Gradle 中路径为:playground)下所有可执行任务,包括构建、Docker 编排、代码生成等,是探索任务体系最直接的入口。
3. 重新生成 protobuf 代码
Playground 的前后端通过 gRPC 通信,接口定义位于 playground/api/v1/api.proto。修改.proto文件后,需要重新生成客户端与服务端桩代码:
cd beam ./gradlew playground:generateProto从 playground/build.gradle.kts 的实现可以看到,generateProto任务底层实际调用buf generate(buf是新一代 protobuf 工具链,在仓库根目录的 playground/buf.gen.yaml 中配置了生成规则)。该文件同时还定义了一个lintProto任务,执行buf lint --path api/对 proto 文件进行规范检查,建议在修改接口定义后一并运行。
// playground/build.gradle.kts task("generateProto") { group = "build" doLast { exec { executable("buf") args("generate") } } }需要注意的是,buf是独立于 Gradle 的命令行工具,运行上述任务前需先安装,并安装 Go / Dart 的 gRPC 依赖(详见 playground/README.md 的"Setup development prerequisites"一节)。
二、用 Docker Compose 搭建本地运行环境
Playground 的本地环境完全由 Docker Compose 编排,配置文件是 playground/docker-compose.local.yaml。整套环境包含 Redis 缓存、Cloud Datastore 模拟器、若干云函数风格的辅助服务、Router 网关、四个 SDK 的 Runner 服务以及 Flutter 前端。
1. 仅启动 Router(最快验证方式)
Router 是 Playground 的入口网关,负责接收 gRPC 请求、路由到对应 SDK 的 Runner 并缓存结果。如果只是调试 Router 本身,可以只启动它:
# 启动 cd beam ./gradlew playground:backend:containers:router:dockerComposeLocalUp # 停止 cd beam ./gradlew playground:backend:containers:router:dockerComposeLocalDown2. 启动 Router、Runners 与前端(完整环境)
完整启动需要按以下步骤操作:
- 编辑 playground/frontend/playground_components/lib/src/constants/backend_urls.dart,将后端地址覆盖为 playground/docker-compose.local.yaml 中对应的本地地址(文件中已预留
// Uncomment the following lines to use local backend.的注释行,取消注释即可)。 - 启动:
cd beam ./gradlew playground:dockerComposeLocalUp- 停止:
cd beam ./gradlew playground:dockerComposeLocalDowndockerComposeLocalUp任务在 playground/build.gradle.kts 中声明了 6 个dependsOn:Router、Go、Java、Python、SCIO 五个后端容器镜像的构建任务,以及前端镜像构建任务,随后执行docker-compose -f docker-compose.local.yaml up -d。因此首次运行会自动先构建全部 Docker 镜像。
重要提示:这种运行方式并未被长期维护,可能在某些环境无法正常工作,主要供前端团队针对尚未部署的后端测试复杂功能。完整启动大约需要 30 分钟且资源消耗较大,建议只启用你当前需要的那个 SDK 的 Runner。如果不需要某些 Runner,可以:
- 在 playground/build.gradle.kts 的
dockerComposeLocalUp任务中注释掉对它的dependsOn; - 在 playground/docker-compose.local.yaml 中注释掉对应的 Docker 镜像配置。
从 playground/docker-compose.local.yaml 可以看到完整环境的服务拓扑:
| 服务 | 端口 | 作用 |
|---|---|---|
redis | 6379 | 远程缓存,Runner 与 Router 通过CACHE_TYPE=remote+CACHE_ADDRESS=redis:6379接入 |
datastore | 8081 | Cloud Datastore 模拟器,DATASTORE_EMULATOR_HOST=datastore:8081 |
cleanup_snippets | 8080(内部) | 定时清理过期片段的云函数(FUNCTION_TARGET: cleanupSnippets) |
put_snippet | 8080(内部) | 写入片段的云函数(FUNCTION_TARGET: putSnippet) |
increment_snippet_views | 8080(内部) | 片段浏览量自增的云函数(FUNCTION_TARGET: incrementSnippetViews) |
router | 8082 | Playground 网关,APPLY_MIGRATIONS: "True"启动时自动应用 Datastore 迁移 |
go_runner/java_runner/python_runner/scio_runner | 8084 / 8086 / 8088 / 8090 | 各 SDK 的代码执行服务 |
frontend | 1234→8080 | Flutter Web 前端 |
这些端口号与 playground/README.md 中"Deploy examples"一节给出的本地 Runner 地址完全一致(Go=localhost:8084、Java=localhost:8086、Python=localhost:8088、SCIO=localhost:8090),本地部署示例时直接复用即可。SCIO Runner 额外设置了mem_limit: 2048M和ulimits.rss,用于模拟真实生产环境的资源限制。
关于前端如何定位后端,可进一步参考前端文档 playground/frontend/README.md 中的 "Backend Lookup" 一节。
三、清理 Cloud Datastore 中的代码片段
Playground 的示例(snippets)保存在 Cloud Datastore 中,长时间运行后会产生过期数据,需要定期清理。
1. 批量清理过期片段
cd beam ./gradlew playground:backend:removeUnusedSnippet -DdayDiff={int} -DprojectId={string} -Dnamespace={datastore namespace}参数含义:
dayDiff:整数。片段的"最后访问日期"(last visited date)若小于等于当前日期减去dayDiff天,则视为过期(out of date)并被删除;projectId:Cloud Datastore 所在 GCP 项目 ID;namespace:Datastore 命名空间,用于隔离不同来源的数据(默认命名空间在 playground/infrastructure/config.py 中定义为Playground)。
在本地环境(test项目)中运行该任务的等价操作由cleanup_snippets云函数承载,cron 表达式则配置在 playground/backend/properties.yaml 的removing_unused_snippets_cron字段,用于生产环境的定时清理。
2. 删除指定片段
当已知某个片段的 ID 时,可以精确删除单个片段:
cd beam ./gradlew playground:backend:removeSnippet -DsnippetId={string} -DprojectId={string} -Dnamespace={datastore namespace}removeUnusedSnippet与removeSnippet这两个 Gradle 任务定义在 playground/backend/build.gradle.kts 中,它们分别对接cmd/remove_unused_snippets与cmd/remove_snippet两个 Go 命令(源码位于 playground/backend/cmd),在本地开发时可直接用这两个任务维护 Datastore 数据。
3. 绕过缓存运行后端测试
cd beam ./gradlew playground:backend:testWithoutCache该任务用于在禁用缓存的情况下运行后端测试,适合验证那些受缓存影响的行为(例如运行结果缓存键的过期与命中逻辑)。
四、依赖版本升级指南
Playground 依赖 Beam SDK 本身以及 Go、Python、Java、Flutter 等语言工具链,升级时需要在多处同步修改,下面是 TASKS.md 给出的完整清单。
1. 升级引用的 Beam SDK 版本
需要依次更新以下位置:
- playground/backend/containers/java/Dockerfile 中默认的
BEAM_VERSION值; - playground/backend/containers/java/build.gradle 中的
default_beam_version; - playground/backend/containers/python/build.gradle 中的
default_beam_version; - CI 中
playground_examples_ci_reusable.yml(位于.github/workflows/)的 SDK 版本; - 部署指南 playground/terraform/README.md "Deploy Playground to Kubernetes" 一节中的
-Psdk-tag=参数; - Cloud Build 流水线中的
_BEAM_VERSION变量,涉及 playground/infrastructure/cloudbuild/playground_ci_stable.yaml 与 playground/infrastructure/cloudbuild/playground_cd_stable.yaml 两个文件。
BEAM_VERSION直接决定 Runner 镜像中内置的 Beam SDK 版本,是"Playground 能运行哪个版本的 Beam 管道"的关键变量;CI / CD 两套 Cloud Build 配置分开维护,是为了保证验证流水线与部署流水线可以独立升级。
2. 升级 Go 版本
- 修改 playground/backend/go.mod 中的 Go 版本号;
- 在 playground/backend 目录下依次运行
go mod download与go mod tidy刷新依赖锁文件(playground/backend/go.sum); - 更新 playground/backend/containers/router/Dockerfile 的
BASE_IMAGE参数及 playground/backend/containers/router/build.gradle 中的对应配置; - 更新 playground/backend/containers/python/Dockerfile、playground/backend/containers/java/Dockerfile、playground/backend/containers/scio/Dockerfile 中的
GO_BASE_IMAGE参数,以及各自目录下build.gradle中的对应配置; - 更新 playground/backend/containers/go/build.gradle 中
buildArgs的BASE_IMAGE。
Go 版本同时影响 Router 与各语言 Runner 的基础镜像,因此需要全量同步;对 Go Runner 而言,playground/backend/README.md 还提到需要额外准备PREPARED_MOD_DIR(预先下载好的go.mod/go.sum目录),升级 Go 版本后建议一并重新生成。
3. 升级 Python 版本
只需更新 playground/backend/containers/python/Dockerfile 中的BASE_IMAGE(即 Python 基础镜像标签)即可。该镜像同时承载 Python SDK 的执行环境,升级后应通过 Python 示例跑一遍验证。
4. 升级 Java 版本
- 更新 playground/backend/containers/java/Dockerfile 中两行基础镜像版本:
FROM maven:3.8.6-openjdk-11 as dep与FROM apache/beam_java11_sdk:$BEAM_VERSION; - 更新 playground/backend/containers/scio/Dockerfile 的
BASE_IMAGE及 playground/backend/containers/scio/build.gradle 中的对应配置。
Java 与 Scala(SCIO)共用 JVM 工具链,所以两处需要同时调整。
5. 升级 Flutter 版本
Flutter 版本的同步涉及三类文件,项目声明的期望最低版本与由依赖计算出的实际最低版本是分开维护的:
项目文件(声明期望的最低版本,每个 Flutter 工程一处):
- playground/frontend/playground_components/pubspec.yaml,
flutter: '>=x.x.x'行; - playground/frontend/playground_components_dev/pubspec.yaml,
flutter: '>=x.x.x'行; - playground/frontend/pubspec.yaml,
flutter: '>=x.x.x'行; - Tour of Beam 前端 learning/tour-of-beam/frontend/pubspec.yaml,
flutter: '>=x.x.x'行。
生成文件(由包依赖解析出的实际最低版本,每个可运行的 Flutter 应用一处):
- playground/frontend/pubspec.lock,
flutter: ">=x.x.x"行; - learning/tour-of-beam/frontend/pubspec.lock,
flutter: ">=x.x.x"行。
pubspec.lock不会手工修改,而是在上述目录中运行dart pub get自动重新生成,随后人工核对并合入变更。
CI 相关:
- playground/frontend/Dockerfile,
ARG FLUTTER_VERSION=x.x.x行; - CI Workflow 中
FLUTTER_VERSION: x.x.x行,涉及.github/workflows/build_playground_frontend.yml、.github/workflows/playground_frontend_test.yml、.github/workflows/tour_of_beam_frontend_test.yml三个文件。
五、生产部署
Playground 的生产部署基于 Terraform + Kubernetes,完整步骤见部署指南 playground/terraform/README.md。该指南覆盖基础设施(GCP 项目、Cloud Build、GKE 集群等)的创建,以及通过-Psdk-tag=指定 Beam SDK 版本将整套应用与依赖基础设施构建并部署到 Kubernetes 的流程。注意其中的-Psdk-tag=参数与第四节"升级 Beam SDK 版本"步骤 5 是对应关系,升级 SDK 时必须同步更新部署参数。
六、获取 SCIO 示例
SCIO 是 Spotify 基于 Apache Beam 的 Scala 包装库,Playground 支持的 SCIO 示例并非直接提交在 Beam 仓库中,而是基于上游 SCIO 仓库中的示例生成。获取方式如下:
python fetch_scala_examples.py --output_dir <output_dir>脚本会把所有受支持的 SCIO 示例下载到<output_dir>。从源码 playground/infrastructure/fetch_scala_examples.py 可以看到其工作原理:
- 内置
SCIO_EXAMPLES列表,逐项记录了示例在 Scio 仓库中的文件路径(如scio-examples/src/main/scala/com/spotify/scio/examples/MinimalWordCount.scala)、名称、描述、pipeline_options、default_example、分类、复杂度与标签等元数据; - 从
SCIO_REPOSITORY(raw.githubusercontent.com/spotify/scio的main分支)拉取源文件; serialize_tag_to_yaml()把元数据序列化成beam-playground:YAML 块(包含name、description、multifile、pipeline_options、default_example、context_line、categories、complexity、tags字段);insert_tag_into_source()将该 YAML 块以//注释形式插入源码的package声明之前,并重算context_line(示例主体起始行),生成带元数据注释的.scala文件写入--output_dir。
要新增一个 SCIO 示例,只需在SCIO_EXAMPLES列表中追加一项,提供上游仓库路径、示例名称、描述、运行选项、分类与标签即可。脚本的输出目录随后可作为ci_cd.py --subdirs的输入,进入下面的示例部署环节。
七、手动部署示例(CI/CD 脚本)
将仓库中的示例、Katas、测试批量写入 Playground 后端,由 playground/infrastructure/ci_cd.py 完成。它同时支持两种模式:
CI:验证(verify)所有 Beam 示例 / 测试 / Katas 可正常执行;CD:验证通过后,把示例及其运行输出写入 Cloud Datastore(GCD)。
1. 前置要求
- 一个已部署 Playground 后端的 GCP 项目;
- Python 3.9.x;
- 已登录 GCP(
gcloud默认登录或使用服务账号密钥)。
2. 环境变量
| 环境变量 | 含义 | 默认值 |
|---|---|---|
GOOGLE_CLOUD_PROJECT | Playground 后端所在 GCP 项目 ID | 无 |
BEAM_ROOT_DIR | 搜索 Playground 示例的根目录 | 无 |
SDK_CONFIG | SDK 与默认示例配置文件位置 | 无 |
BEAM_EXAMPLE_CATEGORIES | 示例分类配置文件位置 | 无 |
BEAM_USE_WEBGRPC | 使用 grpc-Web 而非 grpc | 否 |
GRPC_TIMEOUT | gRPC 调用超时 | 10 秒 |
BEAM_CONCURRENCY | 并行执行的示例数 | 10 |
SERVER_ADDRESS | 特定 SDK 的 Runner 服务地址 | 无 |
其中SDK_CONFIG指向 playground/sdks.yaml,该文件定义了每个 SDK 的默认示例(SDK_GO/SDK_JAVA/SDK_SCIO默认MinimalWordCount,SDK_PYTHON默认WordCountWithMetrics);BEAM_EXAMPLE_CATEGORIES指向 playground/categories.yaml,其中列出了 Playground 支持的全部示例分类(Side Input、Multiple Outputs、Testing、Schemas、Batch、Streaming、Combiners、Dataframes、Joins、IO、Metrics、Options、Coders、Stateful Processing、Beam SQL、Filtering、Branching、Flatten、Core Transforms、Windowing、Debugging、Quickstart、Emulated Data Source、Tee 等),脚本会用它校验示例声明的分类是否合法。
3. 命令行参数
usage: ci_cd.py [-h] --step {CI,CD} [--namespace NAMESPACE] --datastore-project DATASTORE_PROJECT --sdk {SDK_JAVA,SDK_GO,SDK_PYTHON,SDK_SCIO} --origin {PG_EXAMPLES,TB_EXAMPLES} --subdirs SUBDIRS [SUBDIRS ...]--step:CI验证所有示例/测试/Katas;CD额外将示例与运行输出保存到 Cloud Datastore;--namespace:保存数据时使用的 Datastore 命名空间(默认Playground);--datastore-project:保存数据的 Datastore 项目(仅 CD 步骤必需);--sdk:目标 SDK,取值SDK_JAVA/SDK_GO/SDK_PYTHON/SDK_SCIO;--origin:示例来源标记(pg_examples/pg_snippets的 ORIGIN 字段),取值PG_EXAMPLES(主目录示例)或TB_EXAMPLES(Tour of Beam 示例);--subdirs:限定遍历的子目录列表,相对于BEAM_ROOT_DIR。
从 playground/infrastructure/ci_cd.py 源码可以看到完整流水线:读取环境变量 → 加载分类 →find_examples遍历目录收集示例 → 校验重复名称与冲突数据集 →Verifier逐个执行示例验证 → 若是 CD 步骤则通过DatastoreClient保存目录与示例到 Cloud Datastore。注意--datastore-project仅在 CD 步骤为必填,CI 步骤可省略。
4. 一键部署所有 SDK 的辅助脚本
仓库提供了遍历 Go / Java / Python / SCIO 四种 SDK 的部署脚本模板:
cd playground/infrastructure export BEAM_ROOT_DIR="../../" export SDK_CONFIG="../../playground/sdks.yaml" export BEAM_EXAMPLE_CATEGORIES="../categories.yaml" export BEAM_USE_WEBGRPC=yes export BEAM_CONCURRENCY=4 export PLAYGROUND_DNS_NAME="your registered dns name for Playground" for sdk in go java python scio; do export SDK=$sdk && export SERVER_ADDRESS=https://${SDK}.$PLAYGROUND_DNS_NAME && python3 ci_cd.py --datastore-project $GOOGLE_CLOUD_PROJECT \ --step CD --sdk SDK_${SDK^^} \ --origin PG_EXAMPLES \ --subdirs ./learning/katas ./examples ./sdks done脚本要点:
- 每个 SDK 的 Runner 地址形如
https://<sdk>.<PLAYGROUND_DNS_NAME>,即按子域名区分四种 SDK 的服务; SDK_${SDK^^}利用 Shell 大写转换生成SDK_GO/SDK_JAVA/SDK_PYTHON/SDK_SCIO枚举值;--subdirs覆盖三类示例来源:learning/katas(Kata 练习)、examples(Java 示例)、sdks(Go / Python SDK 自带示例);- 本地部署时无需 DNS 域名,
SERVER_ADDRESS直接使用localhost:8084/8086/8088/8090(对应 playground/README.md 的本地 Runner 端口表),并用DATASTORE_EMULATOR_HOST=localhost:8081指向本地 Datastore 模拟器。
八、维护与贡献的延伸阅读
Playground 是一个跨 Go 后端、Dart/Flutter 前端、Terraform 基础设施的大型工程,除本文的运维任务外,还有以下文档可供深入:
- 后端实现细节与环境变量: playground/backend/README.md、playground/backend/CONTRIBUTE.md;
- 前端实现细节与 Backend Lookup: playground/frontend/README.md、playground/frontend/CONTRIBUTE.md;
- 示例元数据格式与如何添加自有示例: playground/load_your_code.md;
- 基础设施与部署: playground/terraform/README.md。
按照 playground/TASKS.md 的顺序,从playgroundPrecommit开始,逐步完成本地环境搭建、数据清理、依赖升级与示例部署,即可完整掌握 Beam Playground 的日常开发与发布流程。
- 大数据
- 批处理
- 流处理
- 数据工程
【免费下载链接】beam
Apache Beam is a unified programming model for Batch and Streaming data processing.
相关推荐
TEN-framework OpenClaw 示例 Playground 前端本地开发指南:环境搭建、依赖安装与运行调试
TEN framework OpenClaw 示例 Playground 前端本地开发指南:环境搭建、依赖安装与运行调试 TEN framework 仓库中的
人工智能AI Agent多模态语音AI 应用Apache Beam Go SDK 实战指南:示例运行、Dataflow 部署与本地构建测试
Apache Beam Go SDK 实战指南:示例运行、Dataflow 部署与本地构建测试 Apache Beam Go SDK 是 Apache Beam
大数据批处理流处理数据工程Apache Beam Playground 完全指南:本地搭建、示例发布与自定义示例接入
Apache Beam Playground 完全指南:本地搭建、示例发布与自定义示例接入 Beam Playground 是 Apache Beam 官方提供
大数据批处理流处理数据工程
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考