Dataflow 批处理可暂停恢复,并用 Blackwell GPU 加速大模型推理

2026-09-15 22 预计阅读时间: 1 分钟
来源: cloud.google.com AI 摘要 Original link

Disclaimer: This article is an AI-assisted summary. Read it together with the original source when precision matters. The summary may omit context, version differences, or edge cases and is not official documentation.

预计阅读时间:10 分钟

长时间运行的数据流水线最怕两件事:任务接近完成时失败,只能从头重算;昂贵的 GPU 被低优先级作业占用,紧急推理任务却排队等待。Dataflow 新增的两项正式可用能力正面解决了这两个问题:批处理作业支持 Pause/Resume,并支持搭载 NVIDIA RTX PRO 6000 Blackwell Server Edition GPU 的 G4 虚拟机。

这不只是一次硬件升级。它改变了批处理作业的恢复方式,也让团队可以把 Dataflow 同时用于数据准备、特征工程和大模型推理,而不必为每个环节维护独立的计算平台。

Pause/Resume 改变了长任务的失败成本

大型批处理可能持续数小时甚至数天。过去,一旦作业失败,已经完成的计算往往无法直接复用,工程师只能重新提交整个任务。结果是重复读取数据、重复执行转换,并再次占用 CPU、GPU 或 TPU。

Pause/Resume 允许 Dataflow 保留作业已经取得的进展,并在之后继续执行。它适合两类场景:

  • 故障恢复:长时间作业失败后,从已有进度恢复,而不是从零开始。
  • 资源抢占与调度:暂停低优先级批处理,把稀缺的 GPU 或 TPU 让给特征工程、在线评估或紧急推理;高优先级任务结束后再恢复原作业。

这项能力尤其适合存在明显优先级差异的 AI 平台。例如,夜间运行的离线嵌入生成任务可以暂停,为白天的模型评估释放加速器;评估结束后,批处理继续运行。

需要注意,Pause 不是任意位置上的进程冻结。恢复效果仍取决于流水线是否能安全保存执行状态、外部系统是否支持幂等写入,以及作业恢复期间输入数据是否发生变化。数据库写入、第三方 API 调用和自定义外部副作用都应单独设计去重机制。

可以这样把暂停与恢复接入运维流程

下面是一份可改造的 Bash 操作模板。运行前需要安装并登录最新版 Google Cloud CLI,然后替换项目、区域和作业 ID。不同 CLI 版本的命令可用性可能不同,因此脚本先检查对应子命令。

#!/usr/bin/env bash
set -euo pipefail

PROJECT_ID='your-project-id'
REGION='us-central1'
JOB_ID='your-dataflow-job-id'
ACTION="${1:-status}"

gcloud config set project "$PROJECT_ID" >/dev/null

case "$ACTION" in
  pause)
    gcloud dataflow jobs pause "$JOB_ID" \
      --region="$REGION"
    ;;
  resume)
    gcloud dataflow jobs resume "$JOB_ID" \
      --region="$REGION"
    ;;
  status)
    gcloud dataflow jobs describe "$JOB_ID" \
      --region="$REGION" \
      --format='yaml(id,name,currentState,currentStateTime)'
    ;;
  *)
    echo "Usage: $0 {pause|resume|status}" >&2
    exit 2
    ;;
esac

保存为 dataflow-control.sh 后执行:

chmod +x dataflow-control.sh
./dataflow-control.sh status
./dataflow-control.sh pause
./dataflow-control.sh resume

如果当前 Cloud CLI 尚未包含 pauseresume 子命令,可以先更新并检查帮助:

gcloud components update
gcloud dataflow jobs --help

生产环境不要只看命令是否返回成功。更稳妥的控制器应轮询 currentState,确认作业真正进入目标状态后再调度 GPU;恢复后还要监控吞吐量、积压、错误率和外部输出是否重复。

G4 与 RTX PRO 6000 把更大的模型带进流水线

Dataflow 现在支持搭载 NVIDIA RTX PRO 6000 Blackwell Server Edition GPU 的 G4 虚拟机。该 GPU 提供 96GB vGPU 显存和 1.6 TB/s 内存带宽,相比此前常用于推理的 NVIDIA L4,为大模型提供了更大的单卡空间和更高的数据吞吐能力。

根据此次能力范围,Dataflow 作业内可以运行 70B+ 参数模型的推理负载。这里的“可以运行”不等于所有 70B 模型都能以任意精度、上下文长度和批大小直接装入单卡。实际显存还会被以下内容占用:

  • 模型权重及量化元数据;
  • KV cache;
  • 推理框架工作区;
  • 输入批次和中间张量;
  • CUDA 上下文及容器内其他进程。

因此,模型参数量只是容量规划的起点。上线前仍要针对精度格式、最大上下文、批大小和并发数做基准测试。

Dataflow 原有的 RunInference、right fitting 和 GPU-enabled autoscaling 可以继续用于模型调用、资源匹配和扩缩容。团队无需自己搭建完整的 GPU 调度层,但仍要明确模型加载方式、容器依赖、区域容量以及冷启动成本。

一个可改造的 GPU 作业启动模板

下面的命令展示了如何为自定义 Apache Beam Python 流水线准备 G4 工作节点。示例假设 pipeline.py 已接受标准 Dataflow Pipeline Options。G4 机型名和加速器标识可能因区域与项目开放情况而变化,先通过查询命令确认可用值。

export PROJECT_ID='your-project-id'
export REGION='us-central1'
export BUCKET='gs://your-dataflow-bucket'

# 先检查目标区域中实际可用的 G4 机型与 GPU 标识。
gcloud compute machine-types list \
  --filter="zone:(${REGION}-*) AND name~'^g4-'" \
  --format='table(name,zone,guestCpus,memoryMb)'

gcloud compute accelerator-types list \
  --filter="zone:(${REGION}-*) AND name~'rtx'" \
  --format='table(name,zone)'

确认名称后,可以这样提交流水线;请把示例变量改成查询到的真实值:

export G4_MACHINE_TYPE='g4-standard-48'
export GPU_TYPE='nvidia-rtx-pro-6000'

python pipeline.py \
  --runner=DataflowRunner \
  --project="$PROJECT_ID" \
  --region="$REGION" \
  --temp_location="$BUCKET/tmp" \
  --staging_location="$BUCKET/staging" \
  --job_name="blackwell-inference-$(date +%Y%m%d-%H%M%S)" \
  --machine_type="$G4_MACHINE_TYPE" \
  --worker_accelerator="type:${GPU_TYPE};count:1;install-nvidia-driver"

这是一份实践模板,而不是对所有区域和 SDK 版本都固定不变的参数承诺。提交前应确认 Dataflow SDK 版本、G4 区域可用性、GPU 配额、驱动安装方式以及自定义容器的 CUDA 兼容性。若使用 Flex Template,还应把对应 Pipeline Options 暴露为模板参数,而不是写死在镜像中。

上线前应做的四项检查

验证恢复语义。 用接近生产规模的数据主动暂停、恢复并注入一次失败,检查作业是否从预期进度继续,以及输出端是否出现重复记录。

建立优先级规则。 明确哪些作业可以被暂停、最长暂停多久、谁有权恢复,以及释放出的 GPU 应分配给哪类工作负载。否则 Pause/Resume 可能变成人工操作瓶颈。

用真实模型做显存基准。 至少记录模型加载显存、峰值显存、首批延迟、稳定吞吐和不同批大小下的成本。96GB 显存很大,但长上下文和高并发仍可能快速消耗 KV cache。

重新核算成本边界。 暂停可以减少重复计算并提高加速器利用率,但不要默认暂停期间所有费用都会消失。状态保存、存储、网络和相关服务仍可能产生费用,应以实际账单和监控数据验证。

Pause/Resume 让 Dataflow 的长批处理更接近可调度、可恢复的生产任务;G4 与 RTX PRO 6000 则扩大了流水线内直接执行大模型推理的上限。最稳妥的采用方式,是先选择一个运行时间长、失败代价高的批任务验证恢复能力,再用一个可量化吞吐和显存的推理作业评估新 GPU。这样既能看到节省的重算成本,也能避免在模型、区域和配额尚未验证前一次性扩大部署。


相关推荐