长时间运行的数据流水线最怕两件事:任务接近完成时失败,只能从头重算;昂贵的 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 尚未包含 pause 或 resume 子命令,可以先更新并检查帮助:
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。这样既能看到节省的重算成本,也能避免在模型、区域和配额尚未验证前一次性扩大部署。