【免费下载链接】beam
Apache Beam is a unified programming model for Batch and Streaming data processing.
导读
本文围绕 Apache Beam 仓库.test-infra/pipelines/infrastructure/04.template/dataflow-to-bigquery这一 Terraform 模块展开,讲解如何将 Beam 内部测试管线ReadDataflowApiWriteBigQuery构建为 Dataflow Flex Template,并通过 Terraform 工作流完成从源码 Shadow Jar 到 GCS 模板文件的自动化打包。读完本文你将掌握:为什么在 Terraform 中要用null_resource+local-exec替代不存在的 Dataflow 模板资源块、Flex Template 构建命令中每个参数的含义与取值来源、以及如何将这套模块复用到自己的 GCP 项目中。
背景:为 Beam 测试管线构建 Dataflow Flex Template
Apache Beam 仓库中的 .test-infra/pipelines 目录存放着支撑测试基础设施的管线,它们仅供内部使用(源码中标注为@Internal)。这些管线依赖 Dataflow Flex Templates 执行——Flex Template 是一种把自定义镜像、启动参数与元数据文件打包到 GCS 上的模板,可以反复启动同一份作业而无需重新编译。
其中的 ReadDataflowApiWriteBigQuery.java 构建并执行一条「读取 Dataflow API、写入 BigQuery」的管线:它从 Eventarc 发布到 Pub/Sub 的 Dataflow 状态变更事件中恢复作业信息,调用 Dataflow 的GetJob、GetJobMetrics、GetJobExecutionDetailsRPC,最后将结果(以及 API 请求错误、转换错误)写入 BigQuery。本 Terraform 模块正是为这条管线生成 Flex Template 的构建层。
为什么用 Terraform?
截至本文档编写时,Google Cloud 的 Terraform Provider 中没有用于创建 Dataflow Template 的资源块。因此,本模块采用null_resource搭配local-execprovisioner 来实现模板构建:
- null_resource:不关联任何真实基础设施资源,仅作为执行逻辑的载体;
- local-exec:在运行
terraform apply的本地机器上执行任意命令。
Terraform 的价值在于其工作流内提供了干净的资源查找能力(data source)。若改用 Gradle 命令或 bash 脚本来提交 Dataflow 作业,往往需要自行维护一堆硬编码的项目 ID、区域、网络名,非常繁琐;而 Terraform 可以让模板构建所需的镜像仓库、服务账号、网络、子网、存储桶等信息全部通过data块动态查询,参数化程度高、可复现性强。
模块文件结构
| 文件 | 作用 |
|---|---|
| template.tf | 核心构建逻辑:Shadow Jar + Flex Template 构建 |
| variables.tf | 声明模块全部输入变量 |
| common.tfvars | 预置的通用变量默认值 |
| apache-beam-testing.tfvars | Apache Beam 测试项目专属变量 |
| data.tf | 查询既有 GCP 资源(Artifact Registry、服务账号、网络、存储桶等) |
| output.tf | 输出生成的模板 GCS 路径 |
| provider.tf | 配置 Google Provider |
| dataflow-template.json | Flex Template 元数据文件(作业参数定义) |
使用方式:Terraform 工作流
本模块遵循 Terraform 工作流约定来应用,假定工作目录位于 .test-infra/pipelines。注意:本模块不使用 state backend,即不维护.tfstate持久化状态(这与 03.io 等资源型模块不同,那里使用了 GCS bucket backend)。
同时注意-var-file标志引用了 common.tfvars,它提供了带倾向性的变量默认值(opinionated variable defaults)。
针对 apache-beam-testing 项目
DIR=infrastructure/04.template/dataflow-to-bigquery terraform -chdir=$DIR init terraform -chdir=$DIR apply -var-file=common.tfvars -var-file=apache-beam-testing.tfvars针对你自己的 GCP 项目
DIR=infrastructure/04.template/dataflow-to-bigquery terraform -chdir=$DIR init terraform -chdir=$DIR apply -var-file=common.tfvars当你使用自己的项目时,需要通过-var或额外的 tfvars 文件覆盖project、storage_bucket_name等变量(例如参考 apache-beam-testing.tfvars 中project = "apache-beam-testing"、storage_bucket_name = "infra-pipelines-19brjqq5"的写法),因为common.tfvars中并未定义这两个变量,它们是必填项(见下文变量表)。
输入变量详解
variables.tf 声明了以下变量:
| 变量 | 类型 | 说明 | 是否必填 |
|---|---|---|---|
project | string | 资源所属的 Google Cloud 项目 ID | 是(common.tfvars未提供) |
artifact_registry_id | string | Artifact Registry 仓库 ID | 否(默认infra-pipelines) |
region | string | 资源所在的 GCP 区域 | 否(默认us-central1) |
dataflow_worker_service_account_id | string | Dataflow Worker 服务账号 ID | 否(默认infra-pipelines-worker) |
network_name_base | string | 网络资源命名基名 | 否(默认infra-pipelines) |
gradle_project | string | 用于构建 Dataflow 模板的 Gradle 项目 | 否(默认:beam-test-infra-pipelines) |
storage_bucket_name | string | 专用于 infra pipelines 的存储桶名称 | 是(common.tfvars未提供) |
template_image_prefix | string | Dataflow 模板的 Artifact Registry 镜像前缀 | 否(默认dataflow-to-bigquery) |
数据源:动态查找既有 GCP 资源
data.tf 通过 Terraform data source 查询先前由 01.setup、02.network、03.io 模块(详见 infrastructure/README.md)预置的资源,无需在 tfvars 中重复填写:
| data source | 查询目标 | 用途 |
|---|---|---|
google_artifact_registry_repository.default | var.artifact_registry_id对应的仓库 | 拼接 Flex Template 镜像地址 |
google_service_account.dataflow_worker | var.dataflow_worker_service_account_id对应的服务账号 | 作为作业运行的服务账号 |
google_project.default | var.project | 项目元数据 |
google_compute_network.default | var.network_name_base | 作业网络 |
google_compute_subnetwork.default | var.region+var.network_name_base | 作业子网 |
google_storage_bucket.default | var.storage_bucket_name | 模板文件落盘位置 |
模板文件最终 GCS 路径由locals推导:
locals { template_file_gcs_path = "gs://${data.google_storage_bucket.default.name}/templates/dataflow-to-bigquery.json" }该路径同时作为 output.tf 的输出值template_file_gcs_path,供terraform output直接读取,方便后续启动作业时引用。
构建流程:两个 null_resource 串联
template.tf 中定义了两个null_resource,用depends_on保证执行顺序:
第一步:构建 Shadow Jar
resource "null_resource" "shadowjar" { triggers = { id = uuid() } provisioner "local-exec" { working_dir = local.beam_root command = "./gradlew ${var.gradle_project}:shadowJar" } }working_dir指向local.beam_root,即${path.module}/../../../../..,也就是仓库根目录;- 命令在仓库根目录执行
./gradlew :beam-test-infra-pipelines:shadowJar,产出带全部依赖的 fat jar; triggers使用uuid()每次 apply 都会变化,从而强制local-exec每次都重新执行(该模块无持久化 state,本地执行始终触发)。
第二步:调用 gcloud 构建 Flex Template
resource "null_resource" "build_template" { triggers = { id = uuid() } depends_on = [null_resource.shadowjar] provisioner "local-exec" { command = <<EOF gcloud dataflow flex-template build \ ${local.template_file_gcs_path} \ --project=${var.project} \ --image-gcr-path=${var.region}-docker.pkg.dev/${var.project}/${data.google_artifact_registry_repository.default.name}/${var.template_image_prefix} \ --metadata-file=${local.beam_root}/.test-infra/pipelines/infrastructure/04.template/dataflow-to-bigquery/dataflow-template.json \ --sdk-language=JAVA \ --flex-template-base-image=JAVA11 \ --jar=${local.beam_root}/.test-infra/pipelines/build/libs/beam-beam-test-infra-pipelines-latest.jar \ --env=FLEX_TEMPLATE_JAVA_MAIN_CLASS=org.apache.beam.testinfra.pipelines.ReadDataflowApiWriteBigQuery \ --network=${data.google_compute_network.default.name} \ --subnetwork=regions/${var.region}/subnetworks/${data.google_compute_subnetwork.default.name} \ --disable-public-ips \ --service-account-email=${data.google_service_account.dataflow_worker.email} \ --enable-streaming-engine EOF } }各参数含义:
gcloud dataflow flex-template build <GCS路径>:构建命令的目标位置,即上一节locals.template_file_gcs_path推导出的gs://<bucket>/templates/dataflow-to-bigquery.json;--project:目标 GCP 项目;--image-gcr-path:模板镜像存放的 Artifact Registry 路径,形如us-central1-docker.pkg.dev/<project>/<repo>/dataflow-to-bigquery(仓库名与镜像前缀均来自 data source 和变量);--metadata-file:Flex Template 元数据 JSON,即 dataflow-template.json;--sdk-language=JAVA/--flex-template-base-image=JAVA11:SDK 语言与基础镜像版本;--jar:上一步 Shadow Jar 产物beam-beam-test-infra-pipelines-latest.jar;--env=FLEX_TEMPLATE_JAVA_MAIN_CLASS=org.apache.beam.testinfra.pipelines.ReadDataflowApiWriteBigQuery:指定模板启动的主类,即本文主角管线;--network/--subnetwork:作业运行的 VPC 网络与子网,来自 data source 查询;--disable-public-ips:禁用公网 IP(私有网络内运行);--service-account-email:Dataflow Worker 服务账号,来自 data source;--enable-streaming-engine:启用 Dataflow Streaming Engine——因为该管线会持续订阅 Pub/Sub 上的 Eventarc 事件,是流式作业。
Flex Template 元数据与作业参数
dataflow-template.json 定义了模板对外暴露的两个必填参数:
| 参数名 | 类型 | 必填 | 帮助文本 |
|---|---|---|---|
subscription | TEXT | 是 | projects/project/subscriptions/subscription,接收 Eventarc 事件的源 Pub/Sub 订阅 |
dataset | TEXT | 是 | project.dataset,写入 Dataflow 作业数据的 BigQuery 数据集 |
这两个参数分别对应管线Options接口中PubsubReadOptions.getSubscription()与BigQueryWriteOptions.getDataset()两个选项(见 ReadDataflowApiWriteBigQuery.java 的Options extends DataflowJobsOptions, PubsubReadOptions, BigQueryWriteOptions)。也就是说,模板使用者只需在启动作业时提供订阅路径与目标数据集,其余运行环境(网络、服务账号、镜像等)都已由 Terraform 固化。
源码佐证:管线执行什么、写入哪些表
了解模板背后实际运行的管线,能帮你判断构建参数是否合理。从 ReadDataflowApiWriteBigQuery.java 的main方法可见其核心流程:
- 读事件:
PubsubIO.readStrings().fromSubscription(...)从 Eventarc 订阅读取 JSON,经EventarcConversions.fromJson()解析为 Dataflow 作业事件; - 筛选作业:
DataflowFilterEventarcJobs只保留「批处理已 DONE」「流式已 CANCELLED / DRAINED」三类作业(JOB_STATE_DONE + JOB_TYPE_BATCH、JOB_STATE_CANCELLED + JOB_TYPE_STREAMING、JOB_STATE_DRAINED + JOB_TYPE_STREAMING),再Flatten合并; - 调 API:
DataflowGetJobs调用JobsV1Beta3.GetJobRPC 获取作业详情,DataflowGetJobMetrics调用GetJobMetrics,DataflowGetJobExecutionDetails调用GetJobExecutionDetails(GetStageExecutionDetails分支当前被注释,留待未来); - 转 Row 并落库:
JobsToRow、WithAppendedDetailsToRow等将 protobuf 结果转换为 BeamRow,最终由 BigQueryWrites.java 写入 BigQuery。
写入侧的表设计也印证了--enable-streaming-engine与时间分区的必要性:BigQueryWrites 定义了jobs、job_metrics、job_execution_details以及errors_jobs_requests、errors_job_metrics_requests、errors_job_execution_details_requests、errors_conversions_*等错误表,全部按HOUR类型按时间字段分区(create_time、job_create_time、observed_time),其中jobs表还额外按type、location聚类。写入配置使用CREATE_IF_NEEDED+WRITE_APPEND+STREAMING_INSERTS(见 BigQueryWrites.java),这意味着表会在首次写入时自动创建,模板参数中的dataset只需指向已存在的数据集。
在整个测试基础设施中的位置
本模块是 .test-infra/pipelines 测试基础设施的「构建层」,与资源层配合使用。参考 infrastructure/README.md,要运行ReadDataflowApiWriteBigQuery管线,推荐的 Terraform 模块应用顺序为:
- 01.setup — 初始化 GCP 项目(服务、IAM、存储、Artifact Registry 等);
- 02.network — 预置网络;
- 03.io/dataflow-to-bigquery — 预置管线读写的资源(Eventarc Workflow、Pub/Sub topic/subscription、GCS bucket、BigQuery 数据集),该模块使用 GCS bucket 作为 state backend(见其 README.md);
- 04.template/dataflow-to-bigquery —本文模块,可选:仅当你想通过 Dataflow Flex Template 运行管线时才需要。
第 4 步的数据源正是读取前几步预置的资源,因此使用本模块前请确保 01~03 已成功 apply,且project、storage_bucket_name等变量与实际资源保持一致。
常见问题与注意事项
- 本模块不持久化 state:每次
apply都会重新构建 Shadow Jar 与 Flex Template(triggers用uuid()强制触发)。若想跳过重建,需自行引入 state 管理; - 必须在仓库根目录外以
-chdir运行:README 中的命令假定工作目录为.test-infra/pipelines,通过terraform -chdir=$DIR切换到模块目录,同时template.tf用相对路径../../../../..回指仓库根目录执行./gradlew,因此请勿在模块目录内单独运行terraform apply; - 模板使用前提:Flex Template 需要 Artifact Registry 仓库、服务账号、网络/子网、存储桶等资源已存在(由 01~03 提供);
--enable-streaming-engine意味着该模板面向流式持续订阅场景; - 适用边界:本模块是 Apache Beam 测试基础设施的内部实现(
@Internal),其硬编码默认值(如us-central1、infra-pipelines)服务于apache-beam-testing项目,复用到其他项目时请以 tfvars 覆盖所有相关变量。
【免费下载链接】beam
Apache Beam is a unified programming model for Batch and Streaming data processing.
相关推荐
基于 Terraform 预置 GCP 基础设施:为 Apache Beam 的 Dataflow API 读取与 BigQuery 写入测试管道构建 Eventarc/Pub/Sub 数据链路
基于 Terraform 预置 GCP 基础设施:为 Apache Beam 的 Dataflow API 读取与 BigQuery 写入测试管道构建 Even
Neon Pageserver 租户时间线迁移指南:基于远程存储的 Timeline Attach/Detach 机制解析
Neon Pageserver 租户时间线迁移指南:基于远程存储的 Timeline Attach/Detach 机制解析 导读 本文聚焦 Neon 项目中 P
Tart 编排实战:Orchard CLI 使用指南——从上下文配置到标签与资源调度
Tart 编排实战:Orchard CLI 使用指南——从上下文配置到标签与资源调度 本指南面向使用 Orchard 编排多个 Tart 主机的开发者,系统讲解
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考