| 名称 | paidf-orchestration-write-dag # 技能名称保留原样 |
| 描述 | 当用户描述自定义的 PAIDF 编排流水线——即特定有序的阶段组合,例如仅增强、仅自动标注、仅检测+描述、或仅图像属性增强(不做完整自动标注)——且现有 airflow/dags/workflows/ 中的 DAG 均未覆盖该组合,并要求创建新的 Kubernetes DAG 时,使用本技能。也用于检查生成的或现有 DAG 的模型/容器/prompt 选择是否与外部规范文档(如 PAIDF 的 launchable.md)匹配。 |
| 版本 | “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 - airflow - dag |
编写 DAG
本技能的功能
从 airflow/dags/shared/task_groups/ 中现有的共享任务组组合出新的仅 Kubernetes 的 Airflow DAG——这些任务组正是 image_attribute_augmentation_dag.py 和 event_video_generation_dag.py 的构建模块。它不会发明新的任务组逻辑。输出为:
airflow/dags/workflows/<dag_name>_dag/configs/<dag_name>_k8s_manifest.yaml— 计算和组件清单airflow/dags/workflows/<dag_name>_dag/<dag_name>_dag.py— DAG 构建器airflow/dags/workflows/<dag_name>_dag/models/payload.py— payload pydantic 模型airflow/dags/workflows/<dag_name>_dag/callables/— 必需的工作流本地粘合代码,而非可选项。这里使用的每个*TaskGroup都接受一个或多个无默认值的Callable构造参数(例如InputPreparationTaskGroup(prepare_input_callable=...),CosmosTaskGroup(config_generation_callable=..., output_validation_callable=...),ValidatedOutputTaskGroup(validation_callable=...))。这些 callable 是真正新的、每个 DAG 特有的代码——它们将通用的容器参数构建转换为这个流水线的实际输入列表/配置生成/输出验证。不要跳过编写它们,也不要将“不要发明任务逻辑”当作自由发明它们的许可证:对于您使用的每个任务组,找到使用相同任务组的现有 DAG,阅读其等效的 callable,并对其进行调整——删除不适用的部分(例如,当没有下一阶段时,丢弃output_group_id/下游 XCom 关注点)——而不是从空白页编写。
仅 Kubernetes。 本仓库的 main/当前分支仅检查 K8s 清单(参见 image-attribute-augmentation-workflow 和 event-video-generation-workflow 技能)——无 NVCF、无 OSMO。除非用户明确表示其环境支持这些后端,否则不要提供这些后端;如果支持,则将其视为需要单独调查的新范围,而不是本技能的标志。
不限于单一工作流。 任务组在 IAA 和 EVG 之间共享;自定义 DAG 可以混合来自两个谱系的阶段(例如 CosmosTaskGroup + AutoLabelingTaskGroup 而不带 DetectionAndTrackingTaskGroup,或者 ImageAttributeAugmentationTaskGroup 而不带 CosmosTaskGroup,用于对已增强的裁剪进行仅标注的通行)。
这些 DAG 不会作为永久工件检查到仓库。 与 image_attribute_augmentation_dag 或 event_video_generation_dag 不同,本技能生成的 DAG 是一次性的、用户请求的流水线——它不会添加到 image-attribute-augmentation-workflow、event-video-generation-workflow 或任何其他运行技能的 DAG-ID 表中。因为不会有其他任何东西引导用户部署或运行它,本技能拥有完整的生命周期,而不仅仅是编写——下面的步骤 8 不是可选的后续步骤,它是这个 DAG 的唯一就绪检查/部署/触发/监控发生的地方。
利用组件技能
这里的任务组是围绕两个组件流水线的薄 Airflow 包装器,这两个流水线有各自的技能,独立于此仓库——augmentation(image-edit / Cosmos Transfer / Cosmos Predict / image2video 容器)和 auto-labeling(检测、描述、视觉问答、人物属性搜索容器)。DAG 的 container_args 和生成的配置文件字面上就是这些技能文档记录的 CLI 契约——本技能不重新实现该契约,而是指向它。
当且仅当本次会话的工作区碰巧包含它们的 repo 时,这两个技能才可达——从 sdg-workflow 到它们没有子模块、插件注册或其他持久链接。首先尝试 Skill(skill="augmentation") / Skill(skill="auto-labeling");如果找不到技能,则回退到 references/component-skills-excerpt.md——本技能实际需要的 schema/prompt/故障排除内容的一个小型、冻结的摘录。将其视为过时的快照,而不是真实技能的替代品:如果需要的资料不在其中,请说明并询问,而不是推断。
在编写清单任务条目、配置模板、prompt 或问题库之前,始终调用匹配的组件技能——绝不要凭记忆发明配置 schema、模型默认值或 prompt 措辞:
| 编写… | 调用技能 | 阅读 |
|---|---|---|
CosmosTaskGroup 配置 — image-edit(IAA 风格) |
augmentation |
references/configuration-schema.md(端点列表、适配器),references/image-attribute-augmentation.md(image-edit prompts,MCQ exclude_variables,分布配置) |
CosmosTaskGroup 配置 — image2video / Cosmos Transfer/Predict(EVG 风格) |
augmentation |
references/event-video-gen.md,references/config-decision-tree.md(根据您实际拥有的输入→输出形状选择哪个模型/适配器) |
ImageAttributeAugmentationTaskGroup / EventAndPersonAttributeSearchTaskGroup prompts、问题库 |
auto-labeling |
references/prompt-authoring.md,references/event-and-person-attribute-search.md,references/stages/person-attribute-search.md |
DetectionAndTrackingTaskGroup / CaptioningTaskGroup / VisualQATaskGroup / AutoLabelingTaskGroup 配置 |
auto-labeling |
references/stages/ 下匹配的文件(detection-and-tracking.md、captioning.md、visual-qa.md) |
| 来自任一容器的配置验证或运行时错误 | 匹配技能 | 其 references/troubleshooting.md |
config-decision-tree.md 是每当任务组执行的不是直接的 image-edit(IAA 的情况)时首先要打开的文件——它告诉您 Cosmos Transfer 与 Predict 与 image2video,而不是任务组名称。
如果用户的请求未映射到组件技能中记录的配置/prompt 模式,请说明并询问,而不是发明 prompt 文本或配置键。
可用的任务组
直接从其子模块导入每个任务组(与每个现有 DAG 文件匹配——不要从 task_groups 包的 __init__ 导入,因为那是不完整的再导出,且遗漏了 EventAndPersonAttributeSearchTaskGroup)。
| 任务组 | 模块 | 功能 | 需要的清单组件 |
|---|---|---|---|
ValidatePayloadTaskGroup |
validate_payload |
根据 DAG 的 pydantic 模型验证触发 payload。始终是第一步。 | 无 |
InputPreparationTaskGroup |
input_preparation |
从 payload 的 input_path 列出/规范化输入媒体。始终必需。 |
无 |
ServiceLifecycleTaskGroup |
service_lifecycle |
部署/拆除内部 VLM/LLM/image-edit/image2video pod,并等待就绪;当 external_services: true 时按服务跳过。 |
components.endpoints 条目与传入的 ServiceLifecycleSpecs 匹配 |
CosmosTaskGroup |
cosmos |
通过清单任务运行 augmentation 组件(image-edit 或 Cosmos Transfer/Predict/image2video),在输入上动态映射。 |
components.tasks.augmentation(+ VLM/LLM/image-edit 端点) |
ImageAttributeAugmentationTaskGroup |
image_attribute_augmentation |
对 cosmos 增强的裁剪运行人物属性搜索/描述(IAA 的自动标注阶段)。读取 cosmos 输出 XCom。 | components.tasks.event_and_person_attribute_search |
EventAndPersonAttributeSearchTaskGroup |
event_and_person_attribute_search |
EVG 的等价物——对自动标注的视频输出进行人物/事件属性搜索。读取自动标注 XCom。 | components.tasks.event_and_person_attribute_search |
DetectionAndTrackingTaskGroup |
detection_and_tracking |
对视频/cosmos 输出执行 RFDETR 检测 + BoxMOT 跟踪。 | components.tasks.auto_labeling(或清单中的检测任务键——检查基础清单) |
CaptioningTaskGroup |
captioning |
对检测/跟踪输出执行 VLM 描述。 | components.tasks.auto_labeling |
VisualQATaskGroup |
visual_qa |
VLM 视觉问答;为每个 QA 目的实例化一次,使用不同的 group_id/component_name(EVG 使用它两次:异常 QA 和人物属性 QA)。 |
components.tasks.visual_qa |
AutoLabelingTaskGroup |
auto_labeling |
通用的映射自动标注容器执行(检测+跟踪作为一个组件)——用于工作流不需要拆分的 DetectionAndTrackingTaskGroup/CaptioningTaskGroup 阶段时。 |
components.tasks.auto_labeling |
ValidatedOutputTaskGroup |
validated_output |
运行工作流特定的输出验证 callable。始终是报告前的最后一个处理步骤。 | 无 |
PerformanceReportingTaskGroup |
reporting |
可选的 YAML+HTML 性能报告,由 payload 中的 enable_performance_reporting 控制。挂在 validated_output 上作为终端分支——绝不能在 fail_pipeline/pipeline_success 上游。 |
无 |
使用任何任务组之前,阅读其源文件中的 __init__——构造函数 kwargs(input_xcom_task_id、group_id、component_name、prepare_args_callable、prepare_args_op_kwargs)因组而异,决定清单接线和 XCom 数据流。从已使用该组的 image_attribute_augmentation_dag.py 或 event_video_generation_dag.py 复制调用模式——不要猜测 kwargs,也不要假设每个组都暴露 input_xcom_task_id 参数:ValidatedOutputTaskGroup 就没有——它接受一个普通的 op_kwargs 字典(默认 {"payload": ..., "run_id": ...}),因此将上游 XCom 拉入它意味着传递自定义的 op_kwargs/callable,而不是任务 ID 字符串。
工作流本地任务组也存在,在 shared/task_groups/ 之外。 IAA 的 CosmosPostProcessingTaskGroup 位于 airflow/dags/workflows/image_attribute_augmentation_dag/tasks/ 下,并有条件地接线——image_attribute_augmentation_dag.py 仅当 "cosmos_post_processing" in self.builder.manifest.deployment.components.tasks 时包含它。检查每个参考 DAG 自己的 tasks/ 目录,而不仅仅是 shared/task_groups/,寻找像这样的阶段特定组,并以相同方式决定包含:以清单任务存在为门控,不要硬编码。
构造函数 kwarg 默认值不能保证与您的清单匹配。 例如,CosmosTaskGroup 的 internal_augmentation_pool 默认是 internal_image_edit_service_pool,但 IAA 的清单 profile 实际命名为 iaa_internal_image_edit_service_pool——参考 DAG 显式覆盖了 kwarg。始终将每个池/profile 名称形状的默认值与您在步骤 3 中编写的清单进行对比,而不仅仅与任务组的源对比。
数据流
input_preparation → service_lifecycle(与 input_preparation 并行,由 services_ready 门控)→ 第一个处理任务组(读取 input_preparation 的 XCom)→ … → 最后一个处理任务组 → validated_output → [fail_pipeline, pipeline_success] 和 performance_reporting(并行分支)→ service_lifecycle 关闭(ALL_DONE)。
每个后续处理任务组必须读取其前驱的 XCom(input_xcom_task_id 指向前一组的组限定任务 ID,在该参数存在的地方——见上文),而不是总是 input_preparation——例如在 IAA 中,event_and_person_attribute_search 读取 cosmos_augmentation.validate_outputs,而不是 input_preparation.prepare_input。
步骤 0:确认没有现有 DAG 已覆盖此场景
在每次进入步骤 1 之前运行此检查——包括当本技能被直接调用时,因为没有上游保证已正确执行此检查。不要因为请求已经说“自定义”或“新建”就跳过它;用户的表述不是证据。
测试是关于任务组链,而不是 payload:image_attribute_augmentation_dag_k8s 已经运行 input_preparation → cosmos_augmentation → cosmos_post_processing → event_and_person_attribute_search → generate_augmented_dataset → validated_output,而 event_video_generation_dag_k8s 已经运行其自己的完整链(参见该 DAG 文件以获取确切的任务组)。如果请求的流水线的阶段序列——在根据下表映射到任务组后——与这些链之一完全匹配,则这不是 write-dag 的情况,无论触发调用的是什么:
- 不同的
input_path/output_directory/服务 URL/模型覆盖/num_augmentation/variable_distribution/max_imgs是 payload 差异。两个已检查的 DAG 已经通过其现有 payload schema 参数化了所有这些——这些都不证明需要新 DAG。 - 只有结构性差异——相对于现有链添加、移除或重新排序任务组——才是真正的
write-dag情况。
如果检查与现有 DAG 匹配:停止,不写任何文件,并告诉用户改用 image-attribute-augmentation-workflow 或 event-video-generation-workflow(以匹配的链为准)——附上他们描述的 payload 差异,因为该技能的 payload 已经涵盖它们。如果确实是结构性差异,则进入步骤 1。
步骤 1:收集需求
一次性询问以下所有问题,而不是逐个询问:
- 流水线 — 用户用自己的话描述的阶段的有序序列。根据上表将每个阶段映射到任务组。如果模糊(例如“增强和标注”——是否包括验证/QA 通行?),在写任何内容之前询问。
- DAG 名称 — snake_case 标识符。
- 输入路径、输出目录、要处理的最大项数。
- 服务模式 —
external_services: true(用户提供端点 URL)或false(DAG 在集群内部署 VLM/LLM/image-edit pod)。如果是外部,收集所选任务组需要的每个服务 URL。 - 容器镜像 / 模型 — 对于每个处理任务组,根据组件技能的当前默认值(
augmentation或auto-labeling,按上表)确认容器镜像和服务模型,而不是硬编码值。请用户确认或覆盖。 - Prompts / 问题库 / MCQ 变量 — 如果流水线需要领域特定的 prompts(编辑 prompt 模板、验证问题、人物属性问题库),使用匹配组件技能的
prompt-authoring.md/image-attribute-augmentation.md指导来创作,而不是从头开始。
步骤 2:阅读基础清单和组件技能参考
在编写任何内容之前,完整阅读最接近的现有 K8s 清单——如果流水线包含 CosmosTaskGroup/ImageAttributeAugmentationTaskGroup,则为 airflow/dags/workflows/image_attribute_augmentation_dag/configs/image_attribute_augmentation_k8s_manifest.yaml;如果包含检测/跟踪/描述/视觉 QA,则为 airflow/dags/workflows/event_video_generation_dag/configs/event_video_generation_k8s_manifest.yaml。这些是 deployment.profiles、池名称、container_args、密钥和 GPU 配置的事实来源——逐字复制字段,不要发明新的 profile 形状。
同时,调用步骤 1 中选择的组件技能并阅读引用的文件。这是模型默认值、端点 adapter 值(nim、openai.chat.completions、openai.images.edits)和 prompt/配置 schema 的来源——清单的 container_args 和 DAG 生成的配置文件必须与该技能文档一致。
步骤 3:编写 K8s 清单
路径: airflow/dags/workflows/<dag_name>_dag/configs/<dag_name>_k8s_manifest.yaml
逐字复制步骤 2 中选择的基础清单,然后:
- 裁剪
components.tasks以仅保留所选任务组需要的条目(参见表格的“需要的清单组件”列)。 - 裁剪
components.endpoints以仅保留剩余任务所需的内容。 - 应用步骤 1 中确认的容器镜像 / 模型覆盖——将生成的
container_args(服务模型名称、所需标志,例如vllm-omni serve Qwen/Qwen-Image-Edit-2511的--omni、匹配暴露 API 路径的端点适配器)与组件技能参考以及用户引用的任何外部规范进行交叉检查(参见下面的一致性检查)。 - 不要更改 profiles、池、密钥或连接设置。
步骤 4:编写 DAG Python 文件
路径: airflow/dags/workflows/<dag_name>_dag/<dag_name>_dag.py
遵循与流水线任务组匹配的参考 DAG 的结构——image_attribute_augmentation_dag.py 用于 image-edit/IAA 风格链,event_video_generation_dag.py 用于检测/描述/视觉 QA/EVG 风格链,或当流水线混合两个谱系时按共享任务组选择更接近的匹配。无论如何,它是当前、工作的模式——不要复活旧的 webserver/NVCF 时代模板:
- 一个
<DagName>DAGBuilder类包装ComponentBuilder(manifest_path=...)。 build_dag()实例化ValidatePayloadTaskGroup、InputPreparationTaskGroup、ServiceLifecycleTaskGroup(使用ServiceLifecycleSpec元组,仅匹配此 DAG 使用的端点),然后按顺序选择处理任务组,然后ValidatedOutputTaskGroup和可选的PerformanceReportingTaskGroup。- 完全按照参考 DAG 接线
fail_pipeline/pipeline_success/shutdown——报告是validated_output的终端分支,绝不能在结果任务上游。 - 模块级注册使用清单存在性保护:
_DAG_NAME_K8S_MANIFEST = os.environ.get(
"<DAG_NAME_UPPER>_K8S_MANIFEST_PATH",
<DAG_NAME_CONFIG_DIR> / "<dag_name>_k8s_manifest.yaml",
)
<dag_name>_dag_k8s = _create_<dag_name>_dag_if_manifest_exists(
manifest_path=_DAG_NAME_K8S_MANIFEST,
platform="k8s",
description="...",
tags=["sdg", "kubernetes", "<dag_name>"],
)
缺少清单意味着 DAG 在 Airflow 中静默不存在,而不是损坏——这是当前代码库中唯一的“注册”机制。没有 webserver 数据库、没有 DAG 播种步骤、也没有要添加的 NVCF/OSMO 分支。
路径: airflow/dags/workflows/<dag_name>_dag/callables/
对于所选链中的每个任务组,通过调整最接近的现有 DAG 的等效实现来编写所需的 callable(参见上面的工作流本地任务组,了解用于非共享组(如 CosmosPostProcessingTaskGroup)的兄弟 tasks/ 目录)。这是真实的、新的代码——预计每个 DAG 都不同——但其形状(签名、XCom 拉取/推送模式、通过 AirflowFailException 的错误处理)应镜像参考实现,而不是从头设计。
注意跨 DAG 重用 tasks/ callable 时的 payload-schema 耦合。 IAA 的 generate_image_attribute_augmentation_performance_report 在其自己的 _parse_payload 内部按名称针对 ImageAttributeAugmentationDagPayloadConfig 重新验证模板化 payload——尽管它位于 schema 无关的报告原语旁边,但它并非 schema 无关。在不同的 DAG 的 payload 模型上未修改地重用它会针对错误的 schema 重新验证;它甚至可能不会大声失败,因为外部模型上的 before-validator 可以静默回填您的 payload 从未拥有的字段。任何导入特定 *DagPayloadConfig 的 tasks/ callable 都需要与上述其他 callable 相同的处理:调整您自己的副本指向您自己的 payload 模型,不要跨 DAG 导入原始文件。
路径: airflow/dags/workflows/<dag_name>_dag/models/payload.py
遵循您在步骤 4 中选择的参考 DAG 的 payload 模型模式(ImageAttributeAugmentationDagPayloadConfig 或 EventVideoGenerationDagPayloadConfig):顶层 input_path、output_directory、external_services、service_lifecycle: ServiceLifecycleTaskConfig、每个处理任务组一个嵌套的 *TaskConfig 字段、enable_performance_reporting: bool,以及一个 model_validator(mode="before"),将 output_directory/external_services 传播到嵌套配置中,以及一个 model_validator(mode="after"),当 external_services: true 但缺少必需的服务 URL 时快速失败。
步骤 5:注册部署
当前部署流程中没有数据库播种步骤。部署命令是 make sync-dag(从仓库根目录运行)。现在不要运行它——步骤 8 在就绪检查后通过明确的用户确认来门控此步骤:
make sync-dag # 打包并上传 DAG、插件和配置到 S3
Airflow 在其正常的 DAG 文件夹刷新时拾取新的 DAG 文件和清单——无需额外的注册调用。
步骤 6:验证
- YAML 语法:
python -c "import yaml; yaml.safe_load(open('<manifest-path>'))" - DAG 导入:
cd airflow && python -c "import dags.workflows.<dag_name>_dag.<dag_name>_dag" - 任务链非空,且每个处理任务组的
input_xcom_task_id指向其实际上游(不总是input_preparation)。 - 清单/代码一致性: 代码引用的每个
components.tasks.*/components.endpoints.*键都存在于清单中,反之亦然——没有孤立的清单条目。
一致性检查(针对外部规范,例如 launchable.md)
如果用户引用了描述相同流水线的外部文档(模型名称、所需的服务标志、端点契约、工作流模式),则将生成的清单和配置与它进行显式差异比较,并报告任何不匹配,而不是静默调和:
- 服务的模型名称和版本匹配。
- 所需的服务标志匹配(例如
vllm-omni serve Qwen/Qwen-Image-Edit-2511的--omni;从组件技能确认,不要假设标志名称)。 - 端点暴露的 API 路径与清单的
adapter选择匹配(例如,文档要求/v1/chat/completions需要adapter: openai.chat.completions,而不是openai.images.edits)。 - 文档中的任何命名“模式”(例如端到端与仅增强与仅标注)映射到单独的生成 DAG 或文档化的 payload 开关——明确说明是哪一个,如果当前任务组无法表达文档描述的模式(参见约束),请说明而不是近似。
在报告成功之前修复任何失败。
步骤 7:打印摘要
生成的文件
├── airflow/dags/workflows/<dag_name>_dag/configs/<dag_name>_k8s_manifest.yaml
│ 端点: <...>
│ 任务: <...>
├── airflow/dags/workflows/<dag_name>_dag/models/payload.py
├── airflow/dags/workflows/<dag_name>_dag/callables/ (改编自:每个 callable 的源 DAG)
└── airflow/dags/workflows/<dag_name>_dag/<dag_name>_dag.py
任务链: input_preparation → <...> → validated_output
DAG ID: <dag_name>_dag_k8s
针对 <引用的文档> 的一致性检查: <通过 / 列出不匹配>
GPU 占用 — 外部模式: <N> 个 GPU — <任务>: 1 GPU(k8s_gpu_task),... (如果链中每个任务 pod 都是 CPU profile,则为 0)
GPU 占用 — 内部模式: <N> 个 GPU = 外部模式总数 + <端点>: 每个 <n> 个 GPU,...
要注册: make sync-dag
从清单本身计算两行 GPU 占用——绝不要陈述或暗示外部模式为 0 而不检查。 两个独立来源,只有一个是模式相关的:
- 在任一模式下进行自己的本地模型推理的任务 pod。 对于
components.tasks中的每个条目,将其deployment_profile解析到deployment.profiles:一个k8s_gpu_taskprofile(或其他带有gpu:的 profile)声称每个运行中的 pod 占用 1 个 GPU 无论external_services如何,因为它在 pod 内进行推理,而不是调用端点——一个k8s_cpu_task/augmentation_taskprofile 的则声称无。这正是 IAA 外部模式真正为 0 GPU 的原因(其任务——augmentation、cosmos_post_processing、event_and_person_attribute_search——都是 CPU profile),而 EVG 的外部模式不是(detection_and_tracking/captioning/visual_qa都在k8s_gpu_task上运行,在考虑任何端点之前至少需要 3 个 GPU)。 - 内部部署的端点,即该服务的
external_services: false。每个副本:VLM / LLM / image-edit / Cosmos Transfer 或 Predict = 1 GPU;image2video = 2 个 GPU(gpu_count: 2,host_ipc: true)。external_services: true在这里贡献 0。
内部模式总数 = 来源 1 + 来源 2——任务 pod 成本不会在内部模式中消失,它是加性的。在摘要中报告两个总数;不要让用户从清单重新推导,也不要重用不同流水线形状的数字(IAA 的“外部 = 0”不能推广到包含任何 k8s_gpu_task profile 阶段的流水线)。
打印覆盖新 payload 模型中每个字段的完整示例 payload,包括嵌套任务配置和流水线需要的任何 prompt/配置文件路径。
步骤 8:就绪检查、部署、触发和监控
因为这个 DAG 永远不会被添加到运行技能的检查 DAG-ID 表中(见上面的“这些 DAG 不会作为永久工件检查”),没有其他任何东西会引导用户运行它。
默认:下面第 1 项(就绪检查)在步骤 7 之后立即自动运行——直接继续,不要停止。 它是只读的(kubectl get,curl GET),不触及共享状态,因此无需确认。
第 2 项(make sync-dag)不同:它将 DAG 文件、插件和配置上传到 S3 并部署到共享的 Airflow 集群——这是对共享基础设施的状态更改操作。 向用户显示步骤 7 的摘要加上就绪检查结果,然后请求明确确认再运行它。调用 write-dag 授权生成文件;它本身不授权将它们部署到共享集群,因此不要将初始调用视为对此步骤的持续批准——每次都确认。同样的先确认后执行规则也适用于触发实际运行(步骤 3-4),原因相同,它还消耗真实的 GPU/计算,并产生可计费、消耗资源的工作流执行。
- 就绪检查,镜像
image-attribute-augmentation-workflow的 SKILL.md## Scope部分和references/airflow-direct-api.md中的检查清单——相同的集群、相同的控制器,只有目标dag_id不同:- 建立集群连接:向用户询问凭据文件路径;绝不假设一个,绝不回退到默认位置,绝不记录文件内容。
- 控制器 pod 运行中:
kubectl get pods -n sdg-workflow -l "release=sdg-workflow-controller"。 - Airflow API 可通过 ClusterIP 访问。 ClusterIP 不保证从本技能运行的位置可路由——取决于主机到集群的网络路径,因环境而异。首先测试它(
curl --max-time 8 "$AIRFLOW_URL/api/v2/version");如果超时,回退到make port-forward作为后台作业(绝不要前台——它会阻塞),并在后续此步骤的每个 API 调用中使用http://localhost:8080。不要将 ClusterIP 超时视为路由到orchestration-setup的就绪检查失败——这是路由差异,不是集群问题,只要端口转发回退有效。 - 部署后(下面的步骤 2),确认此 DAG 的特定
dag_id已加载且is_paused: False。make sync-dag后 DAG 尚未出现意味着同步/刷新尚未生效,而不是生成失败;重新检查而不是重新运行步骤 4-7。 - 有空闲槽位的池:
k8s_gpu_1,default_pool,以及清单自己的任务池(例如iaa_internal_image_edit_service_pool——从您在步骤 3 中编写的清单中读取实际池名称,不要假设它们与参考 DAG 的匹配)。 - 计算集群 GPU 容量(空闲 vs 总计,而不仅仅是可分配的——集群是共享的)对比步骤 7 中计算的 GPU 占用,针对 payload 将实际使用的服务模式——而不是其他模式的数字。
- 命名空间中的陈旧失败 pod(报告,未经确认所有权不要清理)。
- 如果任何检查失败,路由到
orchestration-setup技能而不是进一步诊断。
- 部署:用户确认后(参见上面记录的默认值),运行
make sync-dag(步骤 5),然后重新运行 Airflow-API DAG 已加载检查。 - 构建并确认 payload 与用户一起——与
image-attribute-augmentation-workflow相同的不可协商规则:每个必需字段(input_path、output_directory、服务模式,以及所选任务组需要的任何外部服务 URL)都来自用户,绝不发明、从先前运行重用或从检查文件填充。 - 触发:将确认的 payload POST 到此 DAG 的
dagRuns端点(或 CLI 等效物)——参见同一技能中的references/airflow-direct-api.md了解认证 + 请求形状。捕获运行 ID。 - 监控:轮询运行状态直到终态(
success/failed),定期报告进度,而不是整个运行期间保持沉默。 - 完成时:根据 payload 模型的
output_directory布局获取并总结输出(并且,如果设置了enable_performance_reporting,则获取生成的报告)——与image-attribute-augmentation-workflow的结果检索形状相同,应用于此 DAG 自己的输出路径。
示例
应触发本技能的示例提示,映射到 PAIDF 的 physical-ai-image-attribute-augmentation/launchable.md 中记录的模式。相同形状的请求适用于其他流水线(例如 EVG 风格的检测/描述/视觉 QA 子集)——这三个是今天存在于此仓库中的 IAA 范围实例。
augmentation模式 — “创建一个新的 K8s DAG,仅运行 IAA Cosmos image-edit 阶段,之后不进行自动标注/人物属性搜索。” →CosmosTaskGroup(+ 条件CosmosPostProcessingTaskGroup),没有ImageAttributeAugmentationTaskGroup,没有最终数据集合并步骤。auto_labeling模式 — “创建一个新的 K8s DAG,对已增强的裁剪运行人物属性描述/查询生成,没有图像编辑生成步骤。” → 仅ImageAttributeAugmentationTaskGroup,没有CosmosTaskGroup;输入准备直接列出预增强的裁剪,而不是组合原始人物 ID 窗格。- 偏置/自定义生成 — “为仅增强的 IAA DAG 生成偏向红色上衣和运动鞋的 payload,每人 5 个变体。” → 不是新 DAG;针对已生成 DAG 的
cosmos.variable_distribution构建的 payload,从augmentation组件技能的cosmos_config.yamlverification_options中获取有效属性值,而不是发明。
e2e 模式不是 write-dag 情况——image_attribute_augmentation_dag_k8s 已经覆盖它。
约束
- 仅 K8s — 不要添加 NVCF/OSMO 清单或注册路径。
- 不要修改现有的 DAG 文件、任务组或组件技能源。 只编写新 DAG 自己的文件(清单、payload 模型、callables、DAG 构建器,加上组件技能的创作参考要求的任何新 prompt/问题库/配置文件)。
- 不要重新实现或分叉任务组类逻辑。 每个处理步骤都是现有
*TaskGroup类(共享的或适用的工作流本地的)的实例,以现有 DAG 已经调用它的方式调用。传入这些任务组的 callables 是预期的新代码——从最接近的现有 DAG 的等效实现调整而来,而不是从空白页发明(参见步骤 4)。 - 不要发明 prompts、配置键或模型默认值。 从匹配的组件技能(
augmentation或auto-labeling)或基础清单中获取。 - 如果用户的流水线需要尚不存在的任务组(例如超出现有
group_id命名约定支持的第三个VisualQATaskGroup目的,或两个独立的CosmosTaskGroup轮次——它没有group_id参数,每个 DAG 目前只能有一个增强步骤),解释差距并请求澄清,而不是编写变通代码。