示例工程【免费下载链接】python-docs-samplesCode samples used on cloud.google.com项目地址https://gitcode.com/GitHub_Trending/py/python-docs-samples点击查看免费下载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.pyApache Beam 流水线源码负责在 Worker 上探测 GPU 并打印一条消息Dockerfile定义带 CUDA、TensorFlow 与 Beam SDK 的 Worker 镜像build.yamlCloud Build 构建配置用于构建并推送镜像到 Container Registryrun.yamlCloud Build 运行配置用于提交 Dataflow 作业并附加 GPUrequirements.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 项目并导出项目 IDexport PROJECTyour-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_gpusdef 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(fUsing 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_optionalFalse时直接抛出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_inputbeam.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, --taggcr.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引用。DockerfileGPU 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 中锁定的tensorflow2.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.txtapache-beam[gcp]2.74.0 tensorflow2.21.0其中apache-beam[gcp]是 Dataflow Runner 必需的 GCP 扩展包。第 3 层Beam SDK 注入。第 37 行是整个镜像最巧妙的部分COPY --fromapache/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 REGIONus-central1 export GPU_TYPEnvidia-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 的替换变量区定义了五个参数及其默认值变量默认值含义_IMAGEsamples/dataflow/tensorflow-gpu:latestWorker 镜像名实际为gcr.io/$PROJECT_ID/$_IMAGE_JOB_NAME空Dataflow 作业名正式运行需显式赋值_TEMP_LOCATION空GCS 临时目录正式运行需显式赋值_REGIONus-central1作业运行区域_GPU_TYPEnvidia-tesla-t4GPU 型号_GPU_COUNT1每台 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 - --runnerDataflowRunner - --project$PROJECT_ID - --region$_REGION - --job_name$_JOB_NAME - --temp_location$_TEMP_LOCATION - --sdk_container_imagegcr.io/$PROJECT_ID/$_IMAGE - --machine_typen1-standard-4 - --experimentworker_acceleratortype:$_GPU_TYPE;count:$_GPU_COUNT;install-nvidia-driver - --experimentuse_runner_v2 - --experimentno_use_multiple_sdk_containers - --disk_size_gb50这些参数构成了 Dataflow GPU 作业的核心配置语义--runnerDataflowRunner指定运行器将流水线提交到云端--sdk_container_image告诉 Dataflow 使用哪个自定义镜像作为 Worker 容器——这是自建镜像跑作业的枢纽参数--machine_typen1-standard-4选用 n1-standard-4 机型。GPU 加速要求与机型配套n1 系列是 T4 GPU 的常见载体--experimentworker_acceleratortype:$_GPU_TYPE;count:$_GPU_COUNT;install-nvidia-driver这是挂载 GPU 的核心实验参数三段以分号分隔GPU 类型如nvidia-tesla-t4、每台 Worker 的 GPU 数量默认 1、install-nvidia-driver让 Dataflow 自动安装 NVIDIA 驱动--experimentuse_runner_v2启用 Runner V2这是自定义容器 GPU 方案的必要条件--experimentno_use_multiple_sdk_containers禁用多 SDK 容器模式保证流水线与 SDK 都在同一个自建镜像里运行--disk_size_gb50为 Worker 预留 50 GB 启动磁盘满足 CUDA/TensorFlow 镜像解压空间需求。服务账号与日志配置文件尾部第 53-57 行还包含两处容易被忽略但重要的配置options: logging: CLOUD_LOGGING_ONLY serviceAccount: projects/$PROJECT_ID/serviceAccounts/$PROJECT_NUMBER-computedeveloper.gserviceaccount.comlogging: CLOUD_LOGGING_ONLY构建日志只写入 Cloud Logging不写存储桶减少日志 I/OserviceAccount显式指定 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(scopesession) def build_image(utils: Utils) - str: yield from utils.cloud_build_submit( image_nameNAME, configbuild.yaml, substitutions{_IMAGE: f{NAME}:{utils.uuid}}, ) pytest.fixture(scopesession) def run_dataflow_job(utils: Utils, bucket_name: str, build_image: str) - str: yield from utils.cloud_build_submit( configrun.yaml, substitutions{ _JOB_NAME: utils.hyphen_name(NAME), _IMAGE: f{NAME}:{utils.uuid}, _TEMP_LOCATION: fgs://{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 流水线原文档 Whats next 指向了更完整的卫星影像处理示例在本仓库中对应 tensorflow-landsat。同目录下还有 tensorflow-landsat-prime优化版与 pytorch-minimalPyTorch 版。对照阅读可以发现它们的 Dockerfile、build.yaml、run.yaml 结构与本示例高度一致差异主要在于业务代码Landsat 示例包含真实的影像读取、裁剪与模型推理逻辑依赖清单不同框架对应不同的requirements.txt锁定版本资源规格生产示例可能需要更大的--machine_type或更多--disk_size_gb。也就是说掌握了本示例的镜像-作业两段式流程后升级到真实业务只需替换main.py与requirements.txt基础设施骨架可以原样复用——这正是最小示例的设计价值所在。赞分享示例工程【免费下载链接】python-docs-samplesCode samples used on cloud.google.com项目地址https://gitcode.com/GitHub_Trending/py/python-docs-samples点击查看免费下载相关推荐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 模型示例工程上一篇binwalk C API封装legacy项目集成Rust功能方案下一篇告别繁琐操作用autocmd让Neovim自动完成编辑任务创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
阅读完成 · 觉得有帮助?