Buckets:
| # CASTLE:Vertex Batch 与短 HF Jobs | |
| 当前交付是代码和离线测试。尚未配置 GCS / ADC,也没有提交真实 Batch、创建 HF 定时任务或修改 IAM。 | |
| ## 是否符合“提交后关掉脚本,过一会再取” | |
| 是。输入 JSONL、图片和音频先放入 GCS,Vertex 接收 Batch 后独立执行,HF 提交任务即可退出。后续 HF worker 每次检查一次状态;未完成就退出,完成则收集结果并提交下一阶段。 | |
| 这里是两种 HF worker,而非整个流程恰好只有两个 Job 实例: | |
| 1. `submit`:下载一个源视频到本地磁盘 → 按 clip 提取并上传媒体 → 提交 audio Batch → 退出。 | |
| 2. `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 更快。请按实际模型价格核算。[官方能力与限制](https://docs.cloud.google.com/gemini-enterprise-agent-platform/models/capabilities/batch-inference) | |
| ## 需要配置什么 | |
| | 项目 | 用途 | | |
| |---|---| | |
| | 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 输入与权限说明](https://docs.cloud.google.com/gemini-enterprise-agent-platform/models/capabilities/batch-inference/new-job-from-cloud-storage) | |
| 本代码不会创建 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 中加载,不打印内容: | |
| ```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,至少包含: | |
| ```text | |
| 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 | |
| ```powershell | |
| 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 已存在后,生成定时任务命令: | |
| ```powershell | |
| 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。 | |
| 在已登录的本机查看或停止: | |
| ```text | |
| 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 接口与恢复 | |
| ```text | |
| 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 与结果格式](https://docs.cloud.google.com/gemini-enterprise-agent-platform/models/capabilities/batch-inference/new-job-from-cloud-storage) Python SDK 使用 ADC 的 Vertex client 及 `batches.create/get/list`。[SDK 文档](https://googleapis.github.io/python-genai/) | |
| ## 验证边界 | |
| 离线测试包含 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.