☰
Apache Beam 测试基础设施:用 Terraform 构建 Dataflow Flex Template 托管数据流管线(Dataflow API → BigQuery)
2026/9/29 5:50:28 网站建设 项目流程

【免费下载链接】beam

Apache Beam is a unified programming model for Batch and Streaming data processing.

项目地址:https://gitcode.com/gh_mirrors/beam18/beam
点击查看免费下载

导读

本文围绕 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.tfvarsApache Beam 测试项目专属变量
data.tf查询既有 GCP 资源(Artifact Registry、服务账号、网络、存储桶等)
output.tf输出生成的模板 GCS 路径
provider.tf配置 Google Provider
dataflow-template.jsonFlex 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 声明了以下变量:

变量类型说明是否必填
projectstring资源所属的 Google Cloud 项目 ID是(common.tfvars未提供)
artifact_registry_idstringArtifact Registry 仓库 ID否(默认infra-pipelines)
regionstring资源所在的 GCP 区域否(默认us-central1)
dataflow_worker_service_account_idstringDataflow Worker 服务账号 ID否(默认infra-pipelines-worker)
network_name_basestring网络资源命名基名否(默认infra-pipelines)
gradle_projectstring用于构建 Dataflow 模板的 Gradle 项目否(默认:beam-test-infra-pipelines)
storage_bucket_namestring专用于 infra pipelines 的存储桶名称是(common.tfvars未提供)
template_image_prefixstringDataflow 模板的 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.defaultvar.artifact_registry_id对应的仓库拼接 Flex Template 镜像地址
google_service_account.dataflow_workervar.dataflow_worker_service_account_id对应的服务账号作为作业运行的服务账号
google_project.defaultvar.project项目元数据
google_compute_network.defaultvar.network_name_base作业网络
google_compute_subnetwork.defaultvar.region+var.network_name_base作业子网
google_storage_bucket.defaultvar.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 定义了模板对外暴露的两个必填参数:

参数名类型必填帮助文本
subscriptionTEXT是projects/project/subscriptions/subscription,接收 Eventarc 事件的源 Pub/Sub 订阅
datasetTEXT是project.dataset,写入 Dataflow 作业数据的 BigQuery 数据集

这两个参数分别对应管线Options接口中PubsubReadOptions.getSubscription()与BigQueryWriteOptions.getDataset()两个选项(见 ReadDataflowApiWriteBigQuery.java 的Options extends DataflowJobsOptions, PubsubReadOptions, BigQueryWriteOptions)。也就是说,模板使用者只需在启动作业时提供订阅路径与目标数据集,其余运行环境(网络、服务账号、镜像等)都已由 Terraform 固化。

源码佐证:管线执行什么、写入哪些表

了解模板背后实际运行的管线,能帮你判断构建参数是否合理。从 ReadDataflowApiWriteBigQuery.java 的main方法可见其核心流程:

  1. 读事件:PubsubIO.readStrings().fromSubscription(...)从 Eventarc 订阅读取 JSON,经EventarcConversions.fromJson()解析为 Dataflow 作业事件;
  2. 筛选作业:DataflowFilterEventarcJobs只保留「批处理已 DONE」「流式已 CANCELLED / DRAINED」三类作业(JOB_STATE_DONE + JOB_TYPE_BATCH、JOB_STATE_CANCELLED + JOB_TYPE_STREAMING、JOB_STATE_DRAINED + JOB_TYPE_STREAMING),再Flatten合并;
  3. 调 API:DataflowGetJobs调用JobsV1Beta3.GetJobRPC 获取作业详情,DataflowGetJobMetrics调用GetJobMetrics,DataflowGetJobExecutionDetails调用GetJobExecutionDetails(GetStageExecutionDetails分支当前被注释,留待未来);
  4. 转 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 模块应用顺序为:

  1. 01.setup — 初始化 GCP 项目(服务、IAM、存储、Artifact Registry 等);
  2. 02.network — 预置网络;
  3. 03.io/dataflow-to-bigquery — 预置管线读写的资源(Eventarc Workflow、Pub/Sub topic/subscription、GCS bucket、BigQuery 数据集),该模块使用 GCS bucket 作为 state backend(见其 README.md);
  4. 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.

项目地址:https://gitcode.com/gh_mirrors/beam18/beam
点击查看免费下载
上一篇:SillyTavern自动化革命:5个高级脚本技巧解放你的AI对话生产力
下一篇:5分钟搞定:BetterJoy让你的Switch手柄在PC上完美重生

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询