paidf-orchestration-write-dagSkill paidf-orchestration-write-dag

此技能用于根据用户描述的自定义PAIDF编排流水线(例如仅增强、仅自动标注、或混合阶段等),利用现有共享任务组生成新的Kubernetes Airflow DAG,并处理完整的生命周期管理(编写清单、DAG代码、payload模型、callables,以及部署、触发和监控)。同时可用于检查生成的DAG是否与外部规范文档(如launchable.md)一致。关键词:PAIDF, 编排流水线, Kubernetes DAG, Airflow, 任务组, 数据合成工厂, 物理AI, 图像属性增强, 自动标注, 组件技能, 生命周期管理。

数据合成工厂 0 次安装 1 次浏览 更新于 9/6/2026
名称 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.pyevent_video_generation_dag.py 的构建模块。它不会发明新的任务组逻辑。输出为:

  1. airflow/dags/workflows/<dag_name>_dag/configs/<dag_name>_k8s_manifest.yaml — 计算和组件清单
  2. airflow/dags/workflows/<dag_name>_dag/<dag_name>_dag.py — DAG 构建器
  3. airflow/dags/workflows/<dag_name>_dag/models/payload.py — payload pydantic 模型
  4. 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-workflowevent-video-generation-workflow 技能)——无 NVCF、无 OSMO。除非用户明确表示其环境支持这些后端,否则不要提供这些后端;如果支持,则将其视为需要单独调查的新范围,而不是本技能的标志。

不限于单一工作流。 任务组在 IAA 和 EVG 之间共享;自定义 DAG 可以混合来自两个谱系的阶段(例如 CosmosTaskGroup + AutoLabelingTaskGroup 而不带 DetectionAndTrackingTaskGroup,或者 ImageAttributeAugmentationTaskGroup 而不带 CosmosTaskGroup,用于对已增强的裁剪进行仅标注的通行)。

这些 DAG 不会作为永久工件检查到仓库。image_attribute_augmentation_dagevent_video_generation_dag 不同,本技能生成的 DAG 是一次性的、用户请求的流水线——它不会添加到 image-attribute-augmentation-workflowevent-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.mdreferences/config-decision-tree.md(根据您实际拥有的输入→输出形状选择哪个模型/适配器)
ImageAttributeAugmentationTaskGroup / EventAndPersonAttributeSearchTaskGroup prompts、问题库 auto-labeling references/prompt-authoring.mdreferences/event-and-person-attribute-search.mdreferences/stages/person-attribute-search.md
DetectionAndTrackingTaskGroup / CaptioningTaskGroup / VisualQATaskGroup / AutoLabelingTaskGroup 配置 auto-labeling references/stages/ 下匹配的文件(detection-and-tracking.mdcaptioning.mdvisual-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_idgroup_idcomponent_nameprepare_args_callableprepare_args_op_kwargs)因组而异,决定清单接线和 XCom 数据流。从已使用该组的 image_attribute_augmentation_dag.pyevent_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 默认值不能保证与您的清单匹配。 例如,CosmosTaskGroupinternal_augmentation_pool 默认是 internal_image_edit_service_pool,但 IAA 的清单 profile 实际命名为 iaa_internal_image_edit_service_pool——参考 DAG 显式覆盖了 kwarg。始终将每个池/profile 名称形状的默认值与您在步骤 3 中编写的清单进行对比,而不仅仅与任务组的源对比。

数据流

input_preparationservice_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_imgspayload 差异。两个已检查的 DAG 已经通过其现有 payload schema 参数化了所有这些——这些都不证明需要新 DAG。
  • 只有结构性差异——相对于现有链添加、移除或重新排序任务组——才是真正的 write-dag 情况。

如果检查与现有 DAG 匹配:停止,不写任何文件,并告诉用户改用 image-attribute-augmentation-workflowevent-video-generation-workflow(以匹配的链为准)——附上他们描述的 payload 差异,因为该技能的 payload 已经涵盖它们。如果确实是结构性差异,则进入步骤 1。


步骤 1:收集需求

一次性询问以下所有问题,而不是逐个询问:

  • 流水线 — 用户用自己的话描述的阶段的有序序列。根据上表将每个阶段映射到任务组。如果模糊(例如“增强和标注”——是否包括验证/QA 通行?),在写任何内容之前询问。
  • DAG 名称 — snake_case 标识符。
  • 输入路径、输出目录、要处理的最大项数。
  • 服务模式external_services: true(用户提供端点 URL)或 false(DAG 在集群内部署 VLM/LLM/image-edit pod)。如果是外部,收集所选任务组需要的每个服务 URL。
  • 容器镜像 / 模型 — 对于每个处理任务组,根据组件技能的当前默认值augmentationauto-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 值(nimopenai.chat.completionsopenai.images.edits)和 prompt/配置 schema 的来源——清单的 container_args 和 DAG 生成的配置文件必须与该技能文档一致。

步骤 3:编写 K8s 清单

路径: airflow/dags/workflows/<dag_name>_dag/configs/<dag_name>_k8s_manifest.yaml

逐字复制步骤 2 中选择的基础清单,然后:

  1. 裁剪 components.tasks 以仅保留所选任务组需要的条目(参见表格的“需要的清单组件”列)。
  2. 裁剪 components.endpoints 以仅保留剩余任务所需的内容。
  3. 应用步骤 1 中确认的容器镜像 / 模型覆盖——将生成的 container_args(服务模型名称、所需标志,例如 vllm-omni serve Qwen/Qwen-Image-Edit-2511--omni、匹配暴露 API 路径的端点适配器)与组件技能参考以及用户引用的任何外部规范进行交叉检查(参见下面的一致性检查)。
  4. 不要更改 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() 实例化 ValidatePayloadTaskGroupInputPreparationTaskGroupServiceLifecycleTaskGroup(使用 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 从未拥有的字段。任何导入特定 *DagPayloadConfigtasks/ callable 都需要与上述其他 callable 相同的处理:调整您自己的副本指向您自己的 payload 模型,不要跨 DAG 导入原始文件。

路径: airflow/dags/workflows/<dag_name>_dag/models/payload.py

遵循您在步骤 4 中选择的参考 DAG 的 payload 模型模式(ImageAttributeAugmentationDagPayloadConfigEventVideoGenerationDagPayloadConfig):顶层 input_pathoutput_directoryexternal_servicesservice_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:验证

  1. YAML 语法: python -c "import yaml; yaml.safe_load(open('<manifest-path>'))"
  2. DAG 导入: cd airflow && python -c "import dags.workflows.<dag_name>_dag.<dag_name>_dag"
  3. 任务链非空,且每个处理任务组的 input_xcom_task_id 指向其实际上游(不总是 input_preparation)。
  4. 清单/代码一致性: 代码引用的每个 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 而不检查。 两个独立来源,只有一个是模式相关的:

  1. 在任一模式下进行自己的本地模型推理的任务 pod。 对于 components.tasks 中的每个条目,将其 deployment_profile 解析到 deployment.profiles:一个 k8s_gpu_task profile(或其他带有 gpu: 的 profile)声称每个运行中的 pod 占用 1 个 GPU 无论 external_services 如何,因为它在 pod 内进行推理,而不是调用端点——一个 k8s_cpu_task/augmentation_task profile 的则声称无。这正是 IAA 外部模式真正为 0 GPU 的原因(其任务——augmentationcosmos_post_processingevent_and_person_attribute_search——都是 CPU profile),而 EVG 的外部模式不是(detection_and_tracking/captioning/visual_qa 都在 k8s_gpu_task 上运行,在考虑任何端点之前至少需要 3 个 GPU)。
  2. 内部部署的端点,即该服务的 external_services: false。每个副本:VLM / LLM / image-edit / Cosmos Transfer 或 Predict = 1 GPU;image2video = 2 个 GPU(gpu_count: 2host_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 getcurl GET),不触及共享状态,因此无需确认。

第 2 项(make sync-dag)不同:它将 DAG 文件、插件和配置上传到 S3 并部署到共享的 Airflow 集群——这是对共享基础设施的状态更改操作。 向用户显示步骤 7 的摘要加上就绪检查结果,然后请求明确确认再运行它。调用 write-dag 授权生成文件;它本身不授权将它们部署到共享集群,因此不要将初始调用视为对此步骤的持续批准——每次都确认。同样的先确认后执行规则也适用于触发实际运行(步骤 3-4),原因相同,它还消耗真实的 GPU/计算,并产生可计费、消耗资源的工作流执行。

  1. 就绪检查,镜像 image-attribute-augmentation-workflowSKILL.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: Falsemake sync-dag 后 DAG 尚未出现意味着同步/刷新尚未生效,而不是生成失败;重新检查而不是重新运行步骤 4-7。
    • 有空闲槽位的池: k8s_gpu_1default_pool,以及清单自己的任务池(例如 iaa_internal_image_edit_service_pool——从您在步骤 3 中编写的清单中读取实际池名称,不要假设它们与参考 DAG 的匹配)。
    • 计算集群 GPU 容量(空闲 vs 总计,而不仅仅是可分配的——集群是共享的)对比步骤 7 中计算的 GPU 占用,针对 payload 将实际使用的服务模式——而不是其他模式的数字。
    • 命名空间中的陈旧失败 pod(报告,未经确认所有权不要清理)。
    • 如果任何检查失败,路由到 orchestration-setup 技能而不是进一步诊断。
  2. 部署:用户确认后(参见上面记录的默认值),运行 make sync-dag(步骤 5),然后重新运行 Airflow-API DAG 已加载检查。
  3. 构建并确认 payload 与用户一起——与 image-attribute-augmentation-workflow 相同的不可协商规则:每个必需字段(input_pathoutput_directory、服务模式,以及所选任务组需要的任何外部服务 URL)都来自用户,绝不发明、从先前运行重用或从检查文件填充。
  4. 触发:将确认的 payload POST 到此 DAG 的 dagRuns 端点(或 CLI 等效物)——参见同一技能中的 references/airflow-direct-api.md 了解认证 + 请求形状。捕获运行 ID。
  5. 监控:轮询运行状态直到终态(success/failed),定期报告进度,而不是整个运行期间保持沉默。
  6. 完成时:根据 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.yaml verification_options 中获取有效属性值,而不是发明。

e2e 模式不是 write-dag 情况——image_attribute_augmentation_dag_k8s 已经覆盖它。


约束

  • 仅 K8s — 不要添加 NVCF/OSMO 清单或注册路径。
  • 不要修改现有的 DAG 文件、任务组或组件技能源。 只编写新 DAG 自己的文件(清单、payload 模型、callables、DAG 构建器,加上组件技能的创作参考要求的任何新 prompt/问题库/配置文件)。
  • 不要重新实现或分叉任务组类逻辑。 每个处理步骤都是现有 *TaskGroup 类(共享的或适用的工作流本地的)的实例,以现有 DAG 已经调用它的方式调用。传入这些任务组的 callables 是预期的新代码——从最接近的现有 DAG 的等效实现调整而来,而不是从空白页发明(参见步骤 4)。
  • 不要发明 prompts、配置键或模型默认值。 从匹配的组件技能(augmentationauto-labeling)或基础清单中获取。
  • 如果用户的流水线需要尚不存在的任务组(例如超出现有 group_id 命名约定支持的第三个 VisualQATaskGroup 目的,或两个独立的 CosmosTaskGroup 轮次——它没有 group_id 参数,每个 DAG 目前只能有一个增强步骤),解释差距并请求澄清,而不是编写变通代码。