Buckets:
CASTLE:Vertex Batch 与短 HF Jobs
当前交付是代码和离线测试。尚未配置 GCS / ADC,也没有提交真实 Batch、创建 HF 定时任务或修改 IAM。
是否符合“提交后关掉脚本,过一会再取”
是。输入 JSONL、图片和音频先放入 GCS,Vertex 接收 Batch 后独立执行,HF 提交任务即可退出。后续 HF worker 每次检查一次状态;未完成就退出,完成则收集结果并提交下一阶段。
这里是两种 HF worker,而非整个流程恰好只有两个 Job 实例:
submit:下载一个源视频到本地磁盘 → 按 clip 提取并上传媒体 → 提交 audio Batch → 退出。hourly-tick:每小时启动一个短任务,检查 → 收集 → 推进。audio → annotation → 有需要才执行的 review,分别是独立 Batch。没有音轨时跳过 audio,没有 crop 请求时跳过 review。
GCS state.json 是状态依据,HF /output 保存最终标注、日志和状态快照。Job 重启不需要重新处理整段视频。Batch 是独立执行方式,没有 standard/flex 参数,也不使用在线请求的 AIMD。
目前官方页面列出 Gemini 3.8 Flash 的 Batch 支持和通常比在线低 50% 的价格。它没有预设的用户并发配额,但使用共享容量,仍会排队,不能保证比 standard 更快。请按实际模型价格核算。官方能力与限制
需要配置什么
| 项目 | 用途 |
|---|---|
| Google 项目,启用 Vertex AI / Cloud Storage,开通计费 | 提交并执行 Batch |
| 一个已存在的 GCS bucket,独立 run prefix | 媒体、JSONL、Batch 输出、持久状态;例如 gs://YOUR_BUCKET/castle/smoke-v1 |
| Google ADC 身份 | HF worker 读写 GCS、创建/查看/列出 Batch;当前在线用的 GOOGLE_API_KEY 不能替代这里的 ADC |
| Vertex 服务代理的 bucket 权限 | GCP 内部执行者读取输入和媒体、写出结果 |
| HF 登录态及两个持久目录 | 在本机创建 HF Jobs;代码只读,结果可写 |
常见配置是调用身份拥有项目的 Vertex AI User 权限,及目标 bucket 的对象读写/列举权限。状态文件更新需要覆盖已有对象,因此只授予 Object Creator 不足以支持本管线。可由管理员按实际需求配置更窄的自定义权限。Google 服务代理通常为 service-PROJECT_NUMBER@gcp-sa-aiplatform.iam.gserviceaccount.com;跨项目 bucket 尤其要核对其输入读取和输出写入权限。不要把“Batch 不支持自定义执行服务账号”误读成调用方不能用服务账号认证。GCS 输入与权限说明
本代码不会创建 bucket、开启 API 或授予角色。建议先使用同一项目内、位置合适的专用 bucket,模型位置默认 global。HF bucket (hf://...) 不能充当 GCS (gs://...)。
本机 ADC 与 HF ADC 是两回事
本机调试可使用现有 ADC,或设置 GOOGLE_APPLICATION_CREDENTIALS 指向本机 ADC JSON。HF 容器不会继承本机文件或 gcloud 登录态。
随附 HF launcher 接收环境变量 名称(默认 GOOGLE_ADC_JSON),通过 HF secret 注入 JSON;bootstrap_batch.py 在容器临时磁盘写入仅当前用户可读的 ADC 文件,再移除秘密环境变量并运行 worker。它不把密钥放到 argv、代码 bucket、输出 bucket 或日志中。
如果管理员提供可用于 HF 的服务账号 JSON,在 PowerShell 中加载,不打印内容:
$env:GOOGLE_ADC_JSON = Get-Content -LiteralPath 'C:\secure\castle-adc.json' -Raw
命令执行后可移除本机进程中的变量:Remove-Item Env:GOOGLE_ADC_JSON。用户 ADC 的 refresh token 也属于敏感凭据,需按相同方式保护。若使用工作负载身份联合,须另行配置实际可用的外部 token 来源;仅有 external_account JSON 并不会让 HF 自动获得身份,本 launcher 不配置身份提供方。
本机已登录 HF 时,提交和创建定时任务可以使用登录态;当前 worker 不需要专门注入 HF_TOKEN 来轮询 Vertex。挂载通过 HF Jobs 配置处理。若访问受限数据集、或要求 worker 自己管理 HF schedule,则另行配置相应 HF 权限;本版采用本机/agent 暂停 schedule。
发布代码与生成命令
安装本地依赖:python -m pip install -r requirements-batch.txt。以下命令在 pipeline 目录运行,所有 bucket、项目、源文件都需替换。命令默认只生成 argv,不会提交 HF Job。batch_pipeline.py prepare 本身会上传文件,不是 dry-run。
代码发布到不可变的 HF prefix,至少包含:
batch_pipeline.py batch_jobs.py bootstrap_batch.py
requirements.txt requirements-batch.txt
castle_pipeline/*.py prompts/*.md
不要上传 _test/、密钥、ADC 文件、缓存或完整工作区。已有在线 Dockerfile 不包含 Batch 入口;本 launcher 默认使用 Python 基础镜像并安装 Batch 依赖,submit worker 安装 FFmpeg。正式大规模运行可制作包含相同依赖的镜像,减少每小时启动开销。
一次性提交:先限定一个源文件、3 clips
python batch_jobs.py submit `
--code-volume hf://buckets/Ligant/castle-code/BATCH_CODE_VERSION `
--output-volume hf://buckets/Ligant/castle-output/batch-smoke-v1 `
--name castle-batch-submit-smoke `
--project YOUR_PROJECT `
--state-uri gs://YOUR_BUCKET/castle/smoke-v1/state.json `
--credential-secret GOOGLE_ADC_JSON `
-- `
--gcs-prefix gs://YOUR_BUCKET/castle/smoke-v1 `
--source EXACT_DATASET_VIDEO_PATH `
--model gemini-3.8-flash --max-clips 3
确认范围后,在分隔符 -- 之前加 --execute 才会实际启动,随后 detach。记录返回的 HF Job ID。--source 可以多次给出不同文件,--max-clips 是每文件限制;默认总量硬上限为 2,000 clips。一个小时视频约 120 clips。长范围需相应调整 submit timeout,而不是取消内存边界。
每小时检查一次
首次提交成功、GCS state 已存在后,生成定时任务命令:
python batch_jobs.py hourly-tick `
--code-volume hf://buckets/Ligant/castle-code/BATCH_CODE_VERSION `
--output-volume hf://buckets/Ligant/castle-output/batch-smoke-v1 `
--name castle-batch-tick-smoke `
--project YOUR_PROJECT `
--state-uri gs://YOUR_BUCKET/castle/smoke-v1/state.json `
--credential-secret GOOGLE_ADC_JSON
加 --execute 才会建立 schedule。底层为 HF scheduled run hourly --no-concurrency,单次默认 timeout 30 分钟;不会在一个 worker 中循环等待。记录 schedule ID,它与每次执行的 Job ID 不同。用同样参数把 hourly-tick 换成 tick,可启动一次性的检查 Job。
在已登录的本机查看或停止:
hf jobs inspect JOB_ID
hf jobs logs JOB_ID
hf jobs scheduled inspect SCHEDULE_ID
hf jobs scheduled suspend SCHEDULE_ID
完成、部分完成且有错误、暂停或需要人工处理时,batch-summary.json 给出 stop_schedule: true。本版不会自行暂停 HF schedule;agent/操作者应立即执行 suspend,否则每小时仍有短 HF 实例的费用。终态 tick 不再调用模型或创建 Batch。暂停 HF schedule 也不会取消已经在 GCP 执行的 Batch;若需取消,须明确操作对应 Vertex Job。
本机 worker 接口与恢复
python batch_pipeline.py status --project YOUR_PROJECT --state-uri gs://YOUR_BUCKET/castle/smoke-v1/state.json --output-dir out/batch-check --scratch-dir out/batch-scratch
python batch_pipeline.py tick --project YOUR_PROJECT --state-uri gs://YOUR_BUCKET/castle/smoke-v1/state.json --output-dir out/batch-check --scratch-dir out/batch-scratch
python batch_pipeline.py reconcile --project YOUR_PROJECT --state-uri gs://YOUR_BUCKET/castle/smoke-v1/state.json --output-dir out/batch-check --scratch-dir out/batch-scratch
status 只读取保存状态;tick 会查 GCP,并可能提交下一阶段,因此是运行操作。prepare 不加 --submit 时只准备和初始化,之后 start 可提交。所有 worker 都要求明确 output/scratch 路径。
- 状态写入使用 GCS generation 前置条件;请求和媒体不可覆盖。同内容可重用,不同内容报错。独立运行使用独立 prefix,禁止混用 scope。代码 hash、提示词和范围保存到状态中,后续 tick 要求同版本。
- 提交前先持久化意图,关闭 SDK 自动重试 create。若响应丢失,仅用确定的 display name 查询已存在的 Job。显式
reconcile可在延迟可见后重新查找并接回唯一 Job,永远不会重新 create。找不到或找到多个时需检查 Vertex 控制台,不要修改状态 hash 或盲目重跑提交。 - 音频失败的 clip 不进入视觉阶段;单行坏输出不抹掉其他成功项。原始响应、错误、正规化记录和有效阶段结果保留;不自动重复付费失败请求。需要补跑时先根据状态挑选未完成 clip,使用新的 run prefix 和明确的 clip 范围。本版没有自动跨 run 导入旧阶段结果的接口。
- Batch 输出按回显请求中的
CASTLE_BATCH_ID对齐,不能按文件行号匹配。缺失、重复、未知响应都记为问题;只有通过校验并完成所需复核的 clip 才写 final。 - 断电/终止后,收集可以重放且不会重复提交下一阶段。准备阶段在 state 初始化前中断时,重新准备同配置可复用相同 GCS 对象;仍可能需要重新下载源视频。
- 此 Batch final 在
final/<clip_id>.json,保留annotation、usage、review_required等字段;不直接兼容在线run_pipeline.py status的目录扫描,使用 Batch status/summary。来源与片段范围在 final 内完整记录。
日志、资源和成本
HF stdout 与 /output/batch-events.jsonl 输出媒体开始/完成、每 clip 准备、收集时逐行成功/失败及 token 用量、final 完成、阶段和 Batch Job ID。原始 provider status 保存在 GCS raw 响应,日志只打印归类错误码和可用的结构化 provider code。重放收集可能重复打印记录,按 request_id 去重,不能直接把多次日志相加计算费用。失败响应缺少 usage 不表示没有费用。
等待期间 HF 没有常驻进程,因此没有每 10 分钟心跳或逐请求 AIMD 日志。每小时 tick 提供一次阶段状态,逐条 usage 在阶段结束收集时可见;本版不提前消费仍在运行阶段的输出。
内存按 clip 限定:一个 FFmpeg decoder、一次处理一个片段、一次缩放/裁剪一张图片,JSONL 请求流式写盘,结果逐行落盘再逐 clip 校验。按 4K RGB 估算一张未压缩帧约 24 MiB;30 张 JPEG 保存在磁盘,不在内存中堆叠整小时帧。仍需为解码器、Python、SDK 留余量,真实 HF RSS 尚未实测。磁盘至少容纳当前下载的完整 MP4、一个 clip 的媒体及当阶段结果临时文件。
为支持下个短 Job 的原生分辨率 crop review,启用 review 时还会上传 1 fps 的 native JPEG。一个小时对应约 3,600 张 native 图片 + 3,600 张模型尺寸图片,GCS 存储/操作/传输也计费。--no-review 可省去 native 副本,但改变标注能力,不能在已初始化 run 中途切换。代码不自动删除云端媒体;按实验保留要求另行配置生命周期,不要清理仍被 pending Batch 引用的对象。
官方 JSONL 限制为每批 200,000 请求、1 GB;适配器在上传前流式计数检查,本管线还设更小的 clip 范围上限。官方 JSONL 与结果格式 Python SDK 使用 ADC 的 Vertex client 及 batches.create/get/list。SDK 文档
验证边界
离线测试包含 SDK 请求序列化、create 不自动重试、状态写入竞争、模糊提交恢复、分阶段推进、乱序/缺失/部分失败、native crop、重启重放和最终输出。真实 FFmpeg 使用小合成视频测试 50 fps → 1 fps。
2026-09-16 全套测试:184 passed,含原有在线模式回归。测试结果全部位于专用 _test/ 路径;独立代码审查完成,skill 验证通过并同步到已安装版本。
尚未验证:用户项目的模型访问权限、IAM、GCS 位置、实际费用/队列等待时间、HF 到 GCP 的身份认证,以及真实模型的标注质量。配置完成后应先跑一个源的 2–3 clips,再按明确授权扩展。
Xet Storage Details
- Size:
- 12.8 kB
- Xet hash:
- 67dd842a8198a85738bf320f68488457ae96fb51cdbbbf3354e801f6e835da4f
Xet efficiently stores files, intelligently splitting them into unique chunks and accelerating uploads and downloads. More info.