5.3 Kubeflow Pipelines 的 DAG 编排:把训练变成流水线


文档摘要

5.3 Kubeflow Pipelines 的 DAG 编排:把训练变成流水线 单次训练是「跑一个脚本」,但生产级 ML 需要「数据预处理 → 特征生成 → 训练 → 评估 → 注册 → 部署」一连串步骤协作。把这些步骤组织成有向无环图(DAG),让它们可版本、可缓存、可重跑——这就是 Kubeflow Pipelines 的使命。 5.3.1 从单次训练到流水线:为什么要 DAG 回顾第 4 章 4.1 节的 MLOps 成熟度模型:从第 0 级(手工 notebook)到第 1 级(流水线自动化)的关键跨越,就是把训练流程封装成可重复执行的流水线。 一个生产级训练流程通常包含多个步骤: 数据校验:检查输入数据 schema、缺失率、分布。 特征工程:从特征存储拉特征,做变换。

5.3 Kubeflow Pipelines 的 DAG 编排:把训练变成流水线

单次训练是「跑一个脚本」,但生产级 ML 需要「数据预处理 → 特征生成 → 训练 → 评估 → 注册 → 部署」一连串步骤协作。把这些步骤组织成有向无环图(DAG),让它们可版本、可缓存、可重跑——这就是 Kubeflow Pipelines 的使命。

5.3.1 从单次训练到流水线:为什么要 DAG

回顾第 4 章 4.1 节的 MLOps 成熟度模型:从第 0 级(手工 notebook)到第 1 级(流水线自动化)的关键跨越,就是把训练流程封装成可重复执行的流水线

一个生产级训练流程通常包含多个步骤:

  1. 数据校验:检查输入数据 schema、缺失率、分布。
  2. 特征工程:从特征存储拉特征,做变换。
  3. 训练:跑分布式训练(PyTorchJob)。
  4. 评估:用测试集评估,输出指标。
  5. 模型注册:把模型写入模型注册中心。
  6. 通知:通知团队新模型已就绪。

这些步骤有几个特点:

  • 有依赖关系:评估必须等训练完成。
  • 可独立失败:训练失败不应影响数据校验已成功的结果。
  • 可缓存:如果数据没变,特征工程可缓存,不必重跑。
  • 可并行:无依赖的步骤可并行(如多个评估指标可并行算)。

DAG(Directed Acyclic Graph,有向无环图) 是表达这种「有依赖、可并行」流程的天然数据结构。Kubeflow Pipelines(KFP) 就是 K8s 上承载 ML DAG 流水线的事实标准。

5.3.2 Kubeflow Pipelines 的核心抽象

KFP 围绕几个核心抽象构建:

Component(组件)

流水线的一个步骤,封装为可执行的容器化单元。一个 Component 接受输入、产出输出(Artifact)。Component 通常用 Python SDK 定义,KFP 把它打包成容器镜像。

# KFP 组件伪代码(示意) from kfp import dsl @dsl.component def train(data_path: str, lr: float) -> str: # 训练逻辑 return model_path

Pipeline(流水线)

把多个 Component 按依赖关系组织成 DAG。Pipeline 用 Python decorator 定义,KFP 把它编译成可执行的 YAML/IR。

@dsl.pipeline def my_pipeline(data_path: str): prep_task = preprocess(data_path=data_path) train_task = train(data_path=prep_task.output, lr=0.001) eval_task = evaluate(model_path=train_task.output) register_task = register(model_path=train_task.output, metrics=eval_task.output)

Run(运行)

Pipeline 的一次具体执行。每次 Run 有唯一的 run_id,所有步骤的状态、产物、日志都关联到这个 run_id。

Artifact(产物)

步骤之间传递的数据。可以是模型文件、数据集、指标、图表。KFP 把 Artifact 存在对象存储中,步骤间通过 URL 引用传递。

5.3.3 KFP 的工作流程

KFP 的完整工作流是一个「编译 → 提交 → 执行 → 记录」的循环:

关键机制:

  • 编译期:Python Pipeline 代码被编译成中间表示(IR),描述 DAG 结构。
  • 提交期:IR 提交到 KFP Server,生成一次 Run。
  • 执行期:KFP 按 DAG 拓扑顺序创建 K8s Pod 跑每个 Component。
  • 记录期:每步的输入输出、产物、状态记录到 Metadata,形成血缘。

5.3.4 与 Argo Workflows 的关系

KFP 的底层执行引擎历史上经历了演进:

  • KFP v1:基于 Argo Workflows,KFP 把 Pipeline 编译成 Argo Workflow YAML,由 Argo 调度。
  • KFP v2:迁移到独立的 Argo Runtime / 自研引擎,引入更规范的 IR(PipelineSpec)与更丰富的 ML 语义。

Argo Workflows 本身是一个通用的 K8s 工作流引擎,能编排任意容器化 DAG。它与 KFP 的关系:

维度 KFP Argo Workflows
定位 ML 流水线专用 通用工作流
抽象 Component/Pipeline/Artifact Step/DAG/Template
ML 友好 强(Artifact、Metadata、实验集成)
灵活性 中(约束于 ML 范式) 高(任意容器)
适合 ML 训练流水线 通用 CI/CD、批处理

💡 选型经验:做 ML 流水线 → KFP;做通用工作流(如 CI/CD、批处理脚本编排)→ Argo Workflows。两者底层相通,但 KFP 在 ML 抽象(Artifact、实验跟踪)上更友好。

5.3.5 缓存与重跑:流水线的工程红利

KFP 的一大工程红利是缓存(Caching):如果某个 Component 的输入没变,重新跑 Pipeline 时可直接复用上次的输出,不必重算。这对 ML 流水线价值巨大——比如改了训练超参,但数据没变,那么「数据预处理」步骤可缓存,只重跑训练。

缓存的工作原理:

  • 每个 Component 执行时,KFP 计算其输入的哈希(代码 + 参数 + 上游 Artifact 哈希)。
  • 重新执行时,先查缓存:若哈希命中,直接复用输出。
  • 否则执行 Component,结果入缓存。
缓存命中条件: Component 代码 + 参数 + 上游产物 → 哈希相同 → 命中
场景 缓存效果
改超参,数据没变 数据预处理步骤缓存命中,只重跑训练
改数据,超参没变 全部重跑(数据变了,下游全失效)
改 Pipeline 定义 受影响步骤重跑
完全重跑(debug) 显式指定 cache=False,强制重跑

⚠️ 缓存陷阱:缓存的正确性依赖「哈希能完整反映输入」。如果 Component 内部读了某个外部文件但没在参数中声明,哈希不会变化,缓存会错误命中。这是 KFP 流水线最常见的「为什么我改了东西却没生效」的根因。最佳实践是「所有输入显式声明为参数或 Artifact」。

5.3.6 与训练 Operator 的集成

KFP 的训练 Component 内部通常会启动一个训练作业——这正是第 2 章 2.4 节的 Training Operator 的用武之地。常见的集成模式:

  • KFP Component 提交 PyTorchJob CR:Component 不直接跑训练,而是创建一个 PyTorchJob CR,等其完成后产出模型。
  • PyTorchJob 完成后触发下游:Component 监听 PyTorchJob 状态,完成后把模型路径作为 Artifact 传递给评估 Component。

这种「KFP 编排 + Training Operator 执行训练」的分层,让流水线编排(KFP)与训练执行(Operator)职责分离,是云原生 AI 流水线栈的标准设计。

5.3.7 与第 4 章组件的串联

KFP 流水线把第 4 章的 MLOps 组件串联起来:

  • 数据管理:流水线的输入(DVC/LakeFS 管理的数据版本)。
  • 特征存储:Component 从 Feast 拉特征(4.2 节)。
  • 实验跟踪:每个 Run 自动上报到 MLflow(4.3 节)。
  • 模型注册:流水线最后一个 Component 把模型写入 MLflow Model Registry(4.4 节)。

这样一个 KFP Pipeline 就完成了第 4 章 4.1 节 MLOps 七大组件中的「数据→特征→训练→评估→注册」五步,剩下的「部署→监控」由第 6 章推理服务与第 8 章可观测性承担。

5.3.8 流水线即代码:版本化与可复现

KFP 流水线用 Python 代码定义,意味着:

  • 流水线本身可版本管理:Pipeline 定义在 Git 里,每次修改都有 commit。
  • 可复现:给定 Pipeline 代码 + 输入参数 + 数据版本,能跑出相同结果。
  • 可分享:把 Pipeline 打包成镜像,其他团队可一键运行。
  • 可重跑:发现某个历史 Run 有问题,可用相同参数重跑。

这是 MLOps 1-2 级成熟度的核心特征——「训练流程是代码,可版本可重跑」。对比 0 级的「notebook 手工脚本」,这是质的飞跃。

# 完整 Pipeline 定义伪代码(示意) @dsl.pipeline( name='training-pipeline', description='End-to-end training pipeline' ) def training_pipeline( data_version: str = 'v2.1', lr: float = 0.001, epochs: int = 10, ): # Step 1: 数据校验 validate_task = validate_data(data_version=data_version) # Step 2: 特征拉取 feature_task = fetch_features( data_version=data_version ).after(validate_task) # Step 3: 训练(提交 PyTorchJob) train_task = submit_training( features=feature_task.outputs['features'], lr=lr, epochs=epochs ).after(feature_task) # Step 4: 评估 eval_task = evaluate( model=train_task.outputs['model'] ).after(train_task) # Step 5: 注册(条件分支) with dsl.Condition(eval_task.output > 0.85): register_task = register_model( model=train_task.outputs['model'], metrics=eval_task.output ).after(eval_task)

注意上面用到的 Condition(条件分支):只有评估指标达标时才注册模型。这是 KFP 表达「晋升门禁」的方式,与第 4 章 4.4 节的模型治理结合。

5.3.9 流水线的最佳实践

落地 KFP 流水线的几个工程实践:

  1. 组件细粒度:把流程拆成多个小 Component,而非一个大 Component。细粒度才有缓存红利与失败隔离。
  2. 输入显式化:所有输入(数据、参数、上游 Artifact)显式声明,确保缓存正确。
  3. 幂等性:每个 Component 应可重复执行无副作用(如不直接覆盖、用版本号)。
  4. 错误处理:关键 Component 失败时要有重试与通知机制。
  5. Artifact 治理:定期清理旧 Run 的大 Artifact,平衡追溯与存储成本。
  6. 命名规范:Run 命名包含「项目_模型_日期_关键参数」,便于搜索。
  7. MLflow 集成:每个 Run 自动上报 MLflow,让实验跟踪与流水线 Run 一一对应。
实践 价值
组件细粒度 缓存红利、失败隔离
输入显式化 缓存正确性
幂等性 可重跑
错误处理 鲁棒性
Artifact 治理 成本控制
命名规范 可搜索
MLflow 集成 实验与 Run 对应

💡 判读:KFP 流水线的成熟度,是一个团队 MLOps 1-2 级的关键标志。一个团队如果能用 pipeline.run(data_version='v2.1', lr=0.001) 一行命令重跑任何历史实验,那它已经从 notebook 黑盒跨入了工程化生产。

5.3.10 KFP 之外的其他流水线方案

KFP 不是唯一选择。其他主流方案:

方案 定位 特点
Kubeflow Pipelines ML 流水线 K8s 原生、ML 友好
Argo Workflows 通用工作流 灵活、KFP 底层之一
Tekton K8s CI/CD 标准 CI/CD 友好、Knative 系
Metaflow Netflix 开源 数据科学友好、Python 优先
Prefect / Dagster 现代数据流水线 数据工程友好
Airflow 经典数据流水线 生态广、批处理强

LLM 时代还有 ZenMLFlyte 等新兴方案。选型核心问题:「你的流水线主要是 ML 训练,还是数据工程,还是 CI/CD?」 ML 训练选 KFP,数据工程选 Airflow/Dagster,CI/CD 选 Tekton。

本节小结

  • 流水线把单次训练扩展为「数据→训练→评估→注册」的多步骤 DAG,是 MLOps 0→1 级跨越的关键。
  • KFP 核心抽象:Component(步骤)、Pipeline(DAG)、Run(执行)、Artifact(产物)。
  • KFP 工作流:Python 定义 → 编译 → 提交 → 调度执行 → 元数据记录。
  • KFP v1 基于 Argo Workflows,v2 用独立引擎;KFP 是 ML 专用,Argo 是通用工作流。
  • 缓存是 KFP 的工程红利:输入未变的 Component 直接复用输出;正确性依赖哈希完整反映输入。
  • 「KFP 编排 + Training Operator 执行训练」是云原生 AI 流水线栈的标准分层设计。
  • KFP 流水线把第 4 章的数据/特征/实验/注册串联成完整闭环。
  • 流水线即代码:Pipeline 定义在 Git,可版本、可复现、可分享、可重跑。
  • 流水线支持 Condition 条件分支,可表达模型晋升门禁。
  • 最佳实践:组件细粒度、输入显式化、幂等、错误处理、MLflow 集成。
  • 其他方案:Argo/Tekton(通用)、Airflow/Dagster(数据)、Metaflow(数据科学)。

下一节《5.4 模型训练的 CI/CD 集成》将讲清如何把训练流水线自动化,达成 MLOps 2 级。


发布者: 作者: 灏天文库 转发
评论区 (0)
U