| 名称 | 物理AI事件视频生成 |
| 描述 | 在Kubernetes上运行PAIDF编排事件视频生成DAG——图像到视频异常生成、自动标注和异常数据集生成。选择用于事件视频生成、异常视频生成、图像到视频合成、Cosmos3 image2video、异常数据集创建、安全/监控SDG,或从种子图像生成人员跌倒、人员攀爬、人员奔跑、打斗、吸烟/电子烟、火/烟、入店行窃视频剪辑的请求。当控制器就绪状态未知时,先运行环境设置。不用于人员裁剪衣物/属性增强(那是image-attribute-augmentation-workflow),也不用于视频风格迁移。 |
| 版本 | “1.0.0” |
| 开源协议 | CC-BY-4.0 AND Apache-2.0 metadata: owner: NVIDIA service: physical-ai-data-factory |
| 版本 | 1.0.0 reviewed: ‘2026-09-02’ |
| 作者 | NVIDIA tags: - physical-ai - paidf-orchestration - event-video-generation - cosmos |
PAIDF编排——事件视频生成
端到端运行事件视频生成DAG:种子图像输入准备、Cosmos3图像到视频异常增强、自动标注(检测与跟踪、字幕生成、异常视觉QA、人员属性视觉QA、人员属性搜索)、异常数据集生成和结果检索。
DAG选择
该工作流从airflow/dags/workflows/event_video_generation_dag/为每个计算平台构建一个DAG:
| 平台 | DAG ID | 清单文件 |
|---|---|---|
| Kubernetes | event_video_generation_dag_k8s |
event_video_generation_k8s_manifest.yaml |
Kubernetes是唯一已提交清单文件的平台,因此event_video_generation_dag_k8s是本仓库注册的唯一DAG。只有当清单文件存在时DAG才会被注册;缺失清单意味着DAG在Airflow中不存在,而不是损坏。在触发之前,列出Airflow实际加载的DAG,并且永远不要使用不在该列表中的DAG ID。
只有单一的端到端流水线——没有仅生成或仅标注的DAG变体。如果用户要求只生成视频而不进行自动标注,请告知他们已检入的DAG不提供该流程,而不是编造一个DAG ID。
在Airflow UI中手动输入载荷
如果用户想直接在Airflow UI中输入自己的载荷,而不是让你构建和触发一个,你的工作仅限于让他们进入UI:确认控制器就绪,确保make port-forward正在运行(参见airflow-direct-api.md),并报告可访问的URL。在这种情况下,不要自行渲染载荷、运行预检或触发运行——用户将从UI中完成这些操作。一旦他们告诉你运行已触发,恢复监控(下面第6步);你可以通过Airflow API找到它,而不需要知道他们使用的载荷。
范围
在构建任何载荷之前,从用户那里收集以下所有信息。 不要回退到仓库默认值、CI载荷或任何硬编码的端点URL或存储桶路径。
| 必需 | 字段 | 询问内容 |
|---|---|---|
| 始终 | input_path |
指向单个种子图像或种子图像目录的S3(或HTTP/HTTPS)URL |
| 始终 | output_directory |
结果应写入的可写S3 URL |
| 始终 | 服务模式 | external(用户提供端点URL)或internal(DAG在集群内部署服务) |
| 外部模式 | cosmos.vlm_service_url |
VLM推理端点的完整HTTPS URL |
| 外部模式 | cosmos.llm_service_url |
LLM推理端点的完整HTTPS URL |
| 外部模式 | cosmos.image2video_service_url |
Cosmos3图像到视频推理端点的完整HTTPS URL |
| 可选 | max_images |
从目录处理的图像数量(默认:10;0或负数表示全部) |
| 可选 | cosmos.num_augmentation |
每张图像生成的异常视频数量(默认:1) |
| 可选 | cosmos.variable_distribution |
anomaly_type / env_type采样分布(参见payload-contract.md) |
如果用户未提供必需的值,请在继续之前明确询问。不要编造或复用以前运行或已检入文件中的值。
在触发运行之前,始终运行以下就绪检查。 检查是短路的——在第一次失败时停止,并立即转到环境设置技能。
在任何检查之前,建立集群连接。集群只能通过用户提供的凭据访问——它们从来不是仓库的一部分。检查集群凭据文件路径是否已在shell环境中导出;如果没有,请在运行任何集群命令之前向用户询问绝对路径。永远不要假设路径或回退到任何磁盘上的默认值——完整过程参见setup-and-preflight.md。
控制器(Airflow)和DAG计算任务在同一集群上运行,除非配置了不同的远程集群连接。在此集群上检查GPU容量。
-
控制器Pod——检查Airflow控制器Pod(不是DAG任务Pod)是否处于Running状态。处于
Pending或Failed状态的DAG任务Pod是正常的,不能误认为是控制器故障:kubectl get pods -n sdg-workflow -l "release=sdg-workflow-controller"所有匹配
release=sdg-workflow-controller标签的Pod必须处于Running状态。如果命名空间不存在,这是首次安装的情况——转到环境设置技能,不要进一步诊断。 -
Airflow API——只有在检查1通过后才可以访问。首先从Kubernetes ClusterIP建立
AIRFLOW_URL(始终可从主机路由,无需端口转发):AIRFLOW_URL="http://$(kubectl get svc -n sdg-workflow \ sdg-workflow-controller-api-server \ -o jsonpath='{.spec.clusterIP}'):8080"然后确认目标DAG已加载且
is_paused: False。完整的认证+检查顺序参见airflow-direct-api.md。如果API不可达,转到环境设置技能。 -
池——只有在检查2通过后才进行。具有开放槽位的必需池:
k8s_gpu_1、default_pool,以及所选模式的image2video池(外部模式为external_image2video_service_pool,内部模式为internal_image2video_service_pool)。 -
计算集群GPU——检查集群(使用上面建立的集群连接):
kubectl get nodes \ -o custom-columns='NAME:.metadata.name,GPU_ALLOC:.status.allocatable.nvidia\\.com/gpu' # 同时检查已经消耗GPU的Pod——在共享集群上容量不等于可用性 kubectl get pods -n sdg-workflow \ --field-selector=status.phase=Running -o wide计算集群是共享的——其他用户的运行可能正在进行。将GPU报告为空闲vs总数,而不仅仅是可分配的。两个独立的GPU来源,其中只有一个依赖于模式:
- 无论服务模式如何都运行本地模型推理的任务Pod:
detection_and_tracking、captioning和visual_qa都在k8s_gpu_task配置文件上运行(每个1个GPU)——这个开销在外部和内部模式下都适用,因为这些任务在Pod内进行推理而不是调用端点。event_and_person_attribute_search和augmentation在CPU配置文件上运行,不消耗GPU。因此外部模式至少需要三个GPU,而不是零个。 - 内部部署的端点(
external_services: false):每个VLM/LLM副本一个GPU,加上每个image2video副本两个GPU(gpu_count: 2,host_ipc: true)——每个服务一个副本共四个GPU。 - 内部模式总计 = 两个来源之和:三个任务Pod GPU加上四个端点GPU——至少七个GPU,而不是四个。
- 无论服务模式如何都运行本地模型推理的任务Pod:
-
陈旧的失败Pod——在触发之前,检查计算命名空间中累积的失败Pod并报告它们。它们被设计为保留,不影响运行正确性,但它们消耗命名空间配额并干扰日志搜索:
kubectl get pods -n sdg-workflow \ --field-selector=status.phase=Failed \ -o custom-columns='NAME:.metadata.name,AGE:.metadata.creationTimestamp,DAG:.metadata.labels.dag_id'只有在与用户确认后,才清理
dag_id标签与你拥有的运行匹配的Pod。
明确记录每个检查结果。
如果任何检查失败:自动调用环境设置技能——不要等待用户说“设置”或让他们说出技能名称。
如果用户的请求暗示首次或显式部署(“部署”、“安装”、“设置”、“重新安装”、“重新部署”、“完全设置”):即使所有检查通过,也要调用环境设置技能,并首先确认计划命令。
如果所有检查通过且用户只想运行工作流:直接继续到载荷和触发。
捆绑工具
scripts/upload_images.py: 验证/上传本地种子图像或种子图像的扁平目录。scripts/payload.py: 渲染或验证独立的EventVideoGenerationDagPayloadConfig兼容JSON。scripts/summarize_results.py: 摘要下载的anomaly_dataset/dataset.json数据集。
从此技能目录运行命令。凭据必须从启动代理的shell继承;永远不要要求用户将秘密值粘贴到提示中。
步骤
-
确定输入源。
-
对于本地数据,上传前验证:
python scripts/upload_images.py --path /path/to/seed-images --validate-only -
然后上传:
python scripts/upload_images.py \ --path /path/to/seed-images --destination-path event-video-generation/my-run -
对于现有的存储URL,在确认它指向单个图像或图像的扁平目录(
.jpg、.jpeg、.png、.bmp、.gif、.tiff、.webp)后原样使用。与人员裁剪工作流不同,没有子目录约定——input_path下直接匹配的每个文件都是一个输入图像。当input_path指定一个目录时,DAG对匹配的文件进行排序并取前max_images个。
-
-
选择服务模式。
external需要显式的VLM、LLM和image2video端点URL。internal让DAG的服务生命周期在集群内部署所有三个服务。- 独立于控制器放置选择服务模式。本地控制器可以使用外部推理端点。
- 保持嵌套服务模式和输出目录与顶层一致。
- 在Kubernetes上,VLM和LLM各自从
k8s_gpu_1申请一个GPU,但image2video每个副本申请两个GPU(gpu_count: 2,host_ipc: true)——仅内部部署的端点就需要四个GPU。这是额外的,而不是替代detection_and_tracking/captioning/visual_qa无论服务模式如何都从k8s_gpu_task申请的三个GPU(参见上面的就绪检查GPU分解):外部模式至少需要三个可分配的GPU,内部模式至少需要七个——而不是零个和四个。
-
阅读payload-contract.md,然后根据上面收集的值渲染载荷。不要复制已检入的开发或CI载荷——它们包含部署特定的端点URL和存储桶路径,不得被用户运行继承。
外部模式:
python scripts/payload.py render \ --input-path s3://bucket/input/seed-images/ \ --output-directory s3://bucket/output/event-video-generation/ \ --service-mode external \ --vlm-url https://vlm.example/v1 \ --llm-url https://llm.example/v1 \ --image2video-url https://image2video.example/v1 \ --max-images 10 --num-augmentation 3 \ --variable-distribution assets/variable-distribution.json \ --output /tmp/evg-payload.json内部模式:
python scripts/payload.py render \ --input-path s3://bucket/input/seed-images/ \ --output-directory s3://bucket/output/event-video-generation/ \ --service-mode internal \ --max-images 10 --num-augmentation 3 \ --output /tmp/evg-payload.json向用户展示渲染的载荷(或其验证内容),并在继续之前获得明确确认。只有当他们确认后才继续到预检和触发;如果他们想要更改,重新渲染并重新确认。
-
通过Airflow API对DAG进行预检。检查DAG已加载、必需的池有槽位、控制器Pod健康——参见airflow-direct-api.md#preflight-direct-path。仅确认存在性;永远不要打印凭据值。
-
提交恰好一个DAG运行。将步骤3中的载荷作为
conf.payload传递——完整请求形状参见airflow-direct-api.md#trigger-a-run。记录并返回dag_run_id、输入路径、输出目录和服务模式。 -
触发后立即——无需等待被要求——监控运行直到达到终止状态(
success或failed)。每60-120秒轮询一次Airflow API:# 轮询运行状态 RESPONSE=$(curl -s -H "Authorization: Bearer $TOKEN" \ "$AIRFLOW_URL/api/v2/dags/$DAG_ID/dagRuns/$RUN_ID") RESPONSE="$RESPONSE" python3 -c "import json, os; print(json.loads(os.environ['RESPONSE'])['state'])"当状态为
running或failed时,需要按任务分解的信息,参见airflow-direct-api.md。一旦运行状态为
success或failed就停止轮询。使用适合你运行时的轮询循环——shell的while循环、后台进程或工具原生调度器。不要阻塞用户等待每次轮询;在状态变化发生时报告它们。告诉用户他们也可以在Airflow UI中实时观看进度。
make port-forward在前台运行且从不退出,所以作为后台作业启动——并且最好让用户在自己的终端中运行它,因为代理拥有的转发会随会话终止。解析主机的真实地址,而不是报告占位符或localhost,后者从另一台机器上毫无意义:HOST_IP=$(hostname -I | awk '{print $1}') echo "Airflow UI: http://$HOST_IP:8080"默认凭据是
admin/admin,在deploy/values.yaml中的airflow.createUserJob.defaultUser下定义(不是webserver.defaultUser)。在生产使用前更新它们。完整的按任务分解参见airflow-direct-api.md。
要停止正在进行的运行:打开Airflow UI,找到活动的DagRun,定位正在运行的任务,并将其标记为Failed(任务菜单→标记失败)。这触发DAG的关闭路径,清理Deployments、Services和GPU Pod。不要删除DagRun或DAG——那会绕过清理并留下陈旧的集群资源。
-
在运行达到
success或failed之后,询问用户: “您想下载并分析结果吗?” 不要自动下载——等待确认。如果用户确认,使用shell环境中已有的任何AWS凭据(标准的
AWS_ACCESS_KEY_ID/AWS_SECRET_ACCESS_KEY/AWS_DEFAULT_REGION、AWS配置文件或实例角色)。永远不要要求用户将凭据粘贴到提示中。运行产物位于<output_directory>/<run_id>/下,其中<output_directory>是载荷值,<run_id>是第5步中的dag_run_id。最终数据集在anomaly_dataset/中:aws s3 sync "<output_directory>/<run_id>/anomaly_dataset/" /tmp/evg-results/ python scripts/summarize_results.py --results-dir /tmp/evg-results要检查中间生成的视频,请同步
<output_directory>/<run_id>/cosmos/,并从每个<video_key>/<augmentation_index>/文件夹读取metadata.json。在解释文件之前,请阅读outputs.md。
护栏
- 永远不要使用代码库或已检入载荷中的默认端点URL、存储桶路径或输入路径。 在构建载荷之前,始终向用户询问每个部署特定的值。如果缺少必需的值,停止并询问——不要用猜测代替。
- 在整个会话中保留明确的用户输入和端点/模型选择。
- 如果载荷验证、本地数据集验证或Airflow预检失败,不要提交。
- 不要显示AWS凭据、Airflow bearer令牌或S3签名URL。
- 除非用户明确要求,不要启动多个运行。
- 如果没有提供数据集位置,请询问;此工作流没有隐式演示数据集。
- 不要编造仅生成或仅标注的DAG ID——只有上面列出的DAG存在。
- 只提供清单文件存在且DAG已在Airflow中加载的平台。
- 不要将人员裁剪衣物/属性增强请求路由到这里——那是
image-attribute-augmentation-workflow。
参考资料
- 阅读setup-and-preflight.md了解环境、存储和策略要求。
- 创建或更改载荷时阅读payload-contract.md。
- 检索或解释结果时阅读outputs.md。
- 在验证、API或运行时错误后阅读troubleshooting.md。
- 阅读airflow-direct-api.md了解所有Airflow交互:预检、触发、监控、按任务分解和日志检索。