- 示例工程
【免费下载链接】python-docs-samples
Code samples used on cloud.google.com
Apache Beam 与 Google Cloud Dataflow 为数据并行处理提供了统一模型,而 GPU 的加入让这类流水线可以承担模型推理、图像处理等计算密集任务。本指南以本仓库 tensorflow-minimal 示例为骨架,完整讲解从构建带 TensorFlow 的 Worker 容器镜像,到通过 Cloud Build 在 Dataflow 上拉起 GPU 作业的每一步操作,并结合仓库内源码与配置文件逐行解读其背后的实现原理。读完本文,你将掌握一套可直接复制的"容器化 + GPU 加速" Dataflow 作业搭建方法,并理解worker_accelerator实验参数、Runner V2、多 SDK 容器等关键机制。
示例概览:一个"最小但完整"的 GPU 流水线
tensorflow-minimal位于 dataflow/gpu-examples/tensorflow-minimal,是整个 gpu-examples 目录下最精简的入门示例,它的目录结构本身就是一份部署清单:
- main.py:Apache Beam 流水线源码,负责在 Worker 上探测 GPU 并打印一条消息;
- Dockerfile:定义带 CUDA、TensorFlow 与 Beam SDK 的 Worker 镜像;
- build.yaml:Cloud Build 构建配置,用于构建并推送镜像到 Container Registry;
- run.yaml:Cloud Build 运行配置,用于提交 Dataflow 作业并附加 GPU;
- requirements.txt:流水线依赖锁定文件;
- e2e_test.py 与 noxfile_config.py:端到端测试与测试版本约束。
与目录内其他示例对比(pytorch-minimal 的 PyTorch 版本、tensorflow-landsat 的真实卫星影像处理),该示例刻意把业务逻辑压到最小,让读者把注意力集中在"GPU 基础设施"这一层:镜像怎么构建、作业怎么带 GPU 启动、怎么验证 GPU 真正可用。
前置准备:完成 Dataflow 项目环境初始化
原文档在 "Before you begin" 中要求先完成仓库级 Dataflow setup instructions(对应仓库根目录的 dataflow/README.md)。该指南覆盖以下准备工作:
- 安装 Cloud SDK(在 Cloud Shell 中已预装,可跳过);
- 创建 Google Cloud 项目并导出项目 ID:
export PROJECT=your-google-cloud-project-id; - 使用
gcloud init将 Cloud SDK 绑定到该项目; - 为该 Google Cloud 项目启用结算功能;
- 启用 Dataflow API;
- 执行
gcloud auth application-default login完成本地认证。
需要特别强调的是:本示例在 Docker 镜像内固定了 Python 版本(见下文 Dockerfile 分析),因此本地开发环境的 Python 版本不再是关键约束——这正是该示例"以镜像为运行时"设计思路的一部分。
流水线源码解读:Beam 如何验证 GPU 可用性
main.py 是整个示例唯一的核心业务代码,但它演示了两个值得学习的模式。
首先是 GPU 探测函数 check_gpus:
def check_gpus(_: None, gpus_optional: bool = False) -> None: """Validates that we are detecting GPUs, otherwise raise a RuntimeError.""" gpu_devices = tf.config.list_physical_devices("GPU") if gpu_devices: logging.info(f"Using GPU: {gpu_devices}") elif gpus_optional: logging.warning("No GPUs found, defaulting to CPU.") else: raise RuntimeError("No GPUs found.")它通过tf.config.list_physical_devices("GPU")检查 TensorFlow 是否枚举到 GPU 设备:有则记录日志,没有则在gpus_optional=False时直接抛出RuntimeError,让作业以失败告终——这是一种"fail fast"的校验策略,确保 GPU 环境配置错误能第一时间暴露。
其次,run 函数展示了一个精巧的 Beam 写法:用**旁路输入(side input)**保证 GPU 检查一定执行,且主数据流不受影响:
( pipeline | "Create data" >> beam.Create([input_text]) | "Check GPU availability" >> beam.Map( lambda x, unused_side_input: x, unused_side_input=beam.pvalue.AsSingleton( pipeline | beam.Create([None]) | beam.Map(check_gpus) ), ) | "My transform" >> beam.Map(logging.info) )这里beam.Create([None]) | beam.Map(check_gpus)构成一个独立分支,其结果作为单元素旁路输入喂给主链路的beam.Map。由于 Beam 对旁路输入的处理发生在分布式 Worker 上,check_gpus会在每个 Worker 上实际运行,从而真正做到"在每个 Worker 上探测 GPU"。主元素x原样透传,最终由beam.Map(logging.info)打印到 Worker 日志。
程序入口通过argparse解析--input-text(默认值"Hello!"),剩余参数交给PipelineOptions,这保证了run.yaml中传入的 DataflowRunner 相关参数能被正确识别。
构建 Worker 镜像:Dockerfile 与 build.yaml 详解
原文档指出,镜像构建使用 Cloud Build 并将产物推送到 Container Registry,命令只有一行:
gcloud builds submit --config build.yamlbuild.yaml:镜像构建配置
build.yaml 定义了构建步骤:
substitutions: _IMAGE: samples/dataflow/tensorflow-gpu:latest steps: - name: gcr.io/cloud-builders/docker args: [ build, --tag=gcr.io/$PROJECT_ID/$_IMAGE, . ] images: [ gcr.io/$PROJECT_ID/$_IMAGE ] options: machineType: E2_HIGHCPU_8要点解读:
_IMAGE是用户自定义替换变量,默认推送到gcr.io/$PROJECT_ID/samples/dataflow/tensorflow-gpu:latest;实际部署时(如端到端测试)会覆盖为带随机后缀的镜像名,避免缓存冲突;machineType: E2_HIGHCPU_8指定构建机类型,8 核高 CPU 机型足以应对 TensorFlow 这类体积较大的依赖安装;images字段确保构建产物自动推送到 Container Registry,后续run.yaml可直接以gcr.io/$PROJECT_ID/$_IMAGE引用。
Dockerfile:GPU Worker 镜像的四层结构
Dockerfile 是理解整个方案的关键,它按四层组织:
第 1 层:CUDA 基础镜像。第 18 行:
FROM nvcr.io/nvidia/cuda:12.5.1-cudnn-runtime-ubuntu22.04直接基于 NVIDIA NGC 的 CUDA 12.5.1 + cuDNN runtime 镜像。之所以选-runtime而非-devel,是因为本示例只运行推理级别的 TensorFlow 调用,不需要编译 CUDA 扩展。Dockerfile 注释还提醒读者核对 TensorFlow 与 CUDA 的兼容矩阵(requirements.txt 中锁定的tensorflow==2.21.0与 CUDA 12.5 是配套选型)。
第 2 层:Python 与系统依赖。第 25-33 行通过apt-get安装 Python 3.13,并用update-alternatives将其设为默认python,随后用官方get-pip.py安装 pip,再执行pip install --no-cache-dir -r requirements.txt并做pip check校验依赖完整性。依赖锁定在 requirements.txt:
apache-beam[gcp]==2.74.0 tensorflow==2.21.0其中apache-beam[gcp]是 Dataflow Runner 必需的 GCP 扩展包。
第 3 层:Beam SDK 注入。第 37 行是整个镜像最巧妙的部分:
COPY --from=apache/beam_python3.13_sdk:2.74.0 /opt/apache/beam /opt/apache/beam从官方apache/beam_python3.13_sdk:2.74.0镜像中把 SDK 运行时(/opt/apache/beam)整体拷贝进来,而不是在基础镜像上重新安装。这保证了 Worker 的 Beam SDK 版本与requirements.txt中apache-beam[gcp]==2.74.0严格一致——Dockerfile 注释明确提醒:"Check this matches the apache-beam version in the requirements.txt"。
第 4 层:入口点。第 38 行:
ENTRYPOINT [ "/opt/apache/beam/boot" ]镜像入口固定为 Beam SDK 的boot启动器,这正是 Dataflow Runner V2 容器化架构所要求的形态。
在 Dataflow 上运行带 GPU 的作业:run.yaml 全解
原文档的核心命令如下:
export REGION="us-central1" export GPU_TYPE="nvidia-tesla-t4" gcloud builds submit \ --config run.yaml \ --substitutions _REGION=$REGION,_GPU_TYPE=$GPU_TYPE \ --no-source命令中--no-source意味着不打包任何本地源码——因为代码和依赖早已固化在上一步构建的 Worker 镜像中。Cloud Build 只用run.yaml配置去"调度"一次 Dataflow 作业提交。原文档特别用提示符强调:用 Worker 镜像本身来启动作业,可以保证作业以与 Worker 完全相同的 Python 版本启动,且所有依赖都已就绪。
run.yaml 逐段解析
run.yaml 的替换变量区定义了五个参数及其默认值:
| 变量 | 默认值 | 含义 |
|---|---|---|
_IMAGE | samples/dataflow/tensorflow-gpu:latest | Worker 镜像名(实际为gcr.io/$PROJECT_ID/$_IMAGE) |
_JOB_NAME | 空 | Dataflow 作业名,正式运行需显式赋值 |
_TEMP_LOCATION | 空 | GCS 临时目录,正式运行需显式赋值 |
_REGION | us-central1 | 作业运行区域 |
_GPU_TYPE | nvidia-tesla-t4 | GPU 型号 |
_GPU_COUNT | 1 | 每台 Worker 的 GPU 数量 |
原 README 只导出了REGION和GPU_TYPE两个变量,而 e2e_test.py 的测试夹具补齐了完整视角——它同时替换_JOB_NAME、_IMAGE、_TEMP_LOCATION、_REGION四个变量,说明正式运行时这四个是必填项。
作业提交步骤(第 36-51 行)本质上是用 Worker 镜像中的 Python 直接执行流水线源码:
steps: - name: gcr.io/$PROJECT_ID/$_IMAGE entrypoint: python args: - /pipeline/main.py - --runner=DataflowRunner - --project=$PROJECT_ID - --region=$_REGION - --job_name=$_JOB_NAME - --temp_location=$_TEMP_LOCATION - --sdk_container_image=gcr.io/$PROJECT_ID/$_IMAGE - --machine_type=n1-standard-4 - --experiment=worker_accelerator=type:$_GPU_TYPE;count:$_GPU_COUNT;install-nvidia-driver - --experiment=use_runner_v2 - --experiment=no_use_multiple_sdk_containers - --disk_size_gb=50这些参数构成了 Dataflow GPU 作业的核心配置语义:
--runner=DataflowRunner:指定运行器,将流水线提交到云端;--sdk_container_image:告诉 Dataflow 使用哪个自定义镜像作为 Worker 容器——这是"自建镜像跑作业"的枢纽参数;--machine_type=n1-standard-4:选用 n1-standard-4 机型。GPU 加速要求与机型配套,n1 系列是 T4 GPU 的常见载体;--experiment=worker_accelerator=type:$_GPU_TYPE;count:$_GPU_COUNT;install-nvidia-driver:这是挂载 GPU 的核心实验参数,三段以分号分隔:GPU 类型(如nvidia-tesla-t4)、每台 Worker 的 GPU 数量(默认 1)、install-nvidia-driver让 Dataflow 自动安装 NVIDIA 驱动;--experiment=use_runner_v2:启用 Runner V2,这是自定义容器 + GPU 方案的必要条件;--experiment=no_use_multiple_sdk_containers:禁用多 SDK 容器模式,保证流水线与 SDK 都在同一个自建镜像里运行;--disk_size_gb=50:为 Worker 预留 50 GB 启动磁盘,满足 CUDA/TensorFlow 镜像解压空间需求。
服务账号与日志配置
文件尾部(第 53-57 行)还包含两处容易被忽略但重要的配置:
options: logging: CLOUD_LOGGING_ONLY serviceAccount: projects/$PROJECT_ID/serviceAccounts/$PROJECT_NUMBER-compute@developer.gserviceaccount.comlogging: CLOUD_LOGGING_ONLY:构建日志只写入 Cloud Logging,不写存储桶,减少日志 I/O;serviceAccount:显式指定 Compute Engine 默认服务账号来提交作业,避免使用 Cloud Build 默认权限带来额外授权负担。
验证:日志中应出现 "Using GPU" 输出
由于 main.py 在探测到 GPU 后会输出Using GPU: [...],作业运行后可在 Dataflow 控制台或 Cloud Logging 中检索该关键字确认 GPU 已挂载成功;若配置错误,check_gpus会抛出RuntimeError: No GPUs found.让作业快速失败。
端到端测试:用 nox + Cloud Build 验证整条链路
e2e_test.py 完整复现了"构建镜像 → 提交作业 → 等待完成"的全流程,是理解两个 YAML 如何配合的最佳旁证:
@pytest.fixture(scope="session") def build_image(utils: Utils) -> str: yield from utils.cloud_build_submit( image_name=NAME, config="build.yaml", substitutions={"_IMAGE": f"{NAME}:{utils.uuid}"}, ) @pytest.fixture(scope="session") def run_dataflow_job(utils: Utils, bucket_name: str, build_image: str) -> str: yield from utils.cloud_build_submit( config="run.yaml", substitutions={ "_JOB_NAME": utils.hyphen_name(NAME), "_IMAGE": f"{NAME}:{utils.uuid}", "_TEMP_LOCATION": f"gs://{bucket_name}/temp", "_REGION": utils.region, }, source="--no-source", )测试用utils.uuid为镜像打唯一 tag,随后用同一镜像 tag 提交作业,最后通过utils.dataflow_jobs_wait(job_id)等待作业结束。三个 fixture 均为 session 级,保证镜像只构建一次、作业只提交一次。
noxfile_config.py 则通过ignored_versions跳过 Python 3.8-3.13 的常规矩阵测试,理由是:该示例是 Docker 化样例,Python 版本由 Dockerfile 中的 Beam SDK 容器(apache/beam_python3.13_sdk)决定,跑多个 Python 版本最终都在执行同一个 Dockerfile,没有意义。它还开启了enforce_type_hints: True,与 main.py 中list[str] | None的类型注解风格保持一致。
延伸:从最小示例走向生产级 GPU 流水线
原文档 "What's next" 指向了更完整的卫星影像处理示例,在本仓库中对应 tensorflow-landsat。同目录下还有 tensorflow-landsat-prime(优化版)与 pytorch-minimal(PyTorch 版)。对照阅读可以发现,它们的 Dockerfile、build.yaml、run.yaml 结构与本示例高度一致,差异主要在于:
- 业务代码:Landsat 示例包含真实的影像读取、裁剪与模型推理逻辑;
- 依赖清单:不同框架对应不同的
requirements.txt锁定版本; - 资源规格:生产示例可能需要更大的
--machine_type或更多--disk_size_gb。
也就是说,掌握了本示例的镜像-作业两段式流程后,升级到真实业务只需替换main.py与requirements.txt,基础设施骨架可以原样复用——这正是"最小示例"的设计价值所在。
- 示例工程
【免费下载链接】python-docs-samples
Code samples used on cloud.google.com
相关推荐
python-docs-samples 实战:Dataflow 最小化自定义容器(custom container)从镜像构建到作业运行全流程
python docs samples 实战:Dataflow 最小化自定义容器(custom container)从镜像构建到作业运行全流程 在 Apache
示例工程Dataflow 上运行 PyTorch GPU 最小管道:从镜像构建到作业提交的完整实战
Dataflow 上运行 PyTorch GPU 最小管道:从镜像构建到作业提交的完整实战 导读 本文基于当前仓库 dataflow/gpu examples/
示例工程python-docs-samples 实战:用 Apache Beam RunInference 在 Dataflow 流式管道中运行 Gemma 2B 模型
python docs samples 实战:用 Apache Beam RunInference 在 Dataflow 流式管道中运行 Gemma 2B 模型
示例工程
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考