Volcano JobFlow 作业编排指南:基于 DAG 依赖驱动的云原生批处理工作流引擎
【免费下载链接】volcanoA Cloud Native Batch System (Project under CNCF)项目地址: https://gitcode.com/GitHub_Trending/vol/volcano
Volcano 作为 CNCF 旗下的云原生批量计算系统,以 CRD 方式提供了面向高性能计算与 AI 训练场景的批量作业能力。其中 JobFlow 是 Volcano 针对"多个 VCJob 之间存在先后依赖"这一痛点推出的作业编排引擎,它配合 JobTemplate 模板复用机制,让用户可以像描述 DAG 一样声明作业运行流程。本文以 docs/design/jobflow/README.md 设计文档为核心,结合 pkg/controllers/jobflow 控制器源码与 example/jobflow 示例清单,完整讲解 JobFlow / JobTemplate 的字段语义、依赖判定规则、状态机迁移、Webhook 校验以及端到端使用步骤,帮助你直接在集群中跑通并深度理解其实现原理。
JobFlow 要解决的问题:从"手工编排 VCJob"到"声明式工作流"
Kubernetes 生态中已经有不少工作流引擎,但多数并非为批处理作业设计。批处理作业(如 AI 训练、大数据分析、HPC 任务)通常存在复杂的运行依赖,且单个作业耗时可能长达数天甚至数周。在引入 JobFlow 之前,多个 Volcano VCJob 之间的协作往往需要:
- 人工串行提交、等待、再提交;
- 借助外部作业编排平台手工调度;
- 自行实现"上一个任务完成后触发下一个任务"的轮询逻辑。
JobFlow 的目标正是将这种"作业间依赖"以声明式方式内建到 Volcano 中。它提出两个核心概念(详见 docs/design/jobflow/README.md):
- JobTemplate(作业模板,缩写 jt):VCJob 的模板,定义作业的完整 spec,但不被 job 控制器直接下发,等待被 JobFlow 引用;
- JobFlow(作业流,缩写 jf):定义一组作业的运行流程,通过
flows字段描述作业之间的依赖关系(顺序执行、并行执行、条件依赖等)。
JobFlow 不是通用工作流引擎,它了解 VCJob 的细节,因此能够为用户提供远超通用引擎的作业感知能力,例如:作业运行状态、起止时间戳、下一个要运行的作业、Pod 失败比率等。
从设计文档的 Scope 划分来看:
- In Scope:JobFlow API 与行为定义、多作业之间的启动顺序、作业启动顺序的依赖完成状态、基于 DAG 的作业依赖启动;
- Out of Scope:支持其他作业类型、实现 vcjob 级别的 gang 调度。
也就是说,JobFlow 聚焦"作业与作业之间的编排",而作业内部的 gang 调度、优先级抢占等能力仍由 Volcano 调度器与 VCJob 控制器负责。
整体架构与作业提交流程
JobFlow 属于 Volcano 控制器体系中的一员,在 pkg/controllers/jobflow/jobflow_controller.go 中通过framework.RegisterController(&jobflowcontroller{})注册为jobflow-controller,随 vc-controller-manager 一起运行。
一次完整的 JobFlow 作业提交链路
根据设计文档,一次完整的提交过程如下(可对照架构图 docs/design/images/jobflow-2.png,其中蓝色为 Kubernetes 原生组件、橙色为 Volcano 既有定义、红色为 JobFlow 新增定义):
- 通过 Admission 后,
kubectl在 kube-apiserver 中创建 JobTemplate 与 JobFlow(Volcano CRD)对象; - JobFlowController 以 JobTemplate 为模板,根据 JobFlow 的配置与依赖规则创建对应的 VCJob;
- VCJob 创建后,VCJobController 根据 VCJob 配置创建对应的 Pod 与 PodGroup;
- Pod 与 PodGroup 创建后,vc-scheduler 从 kube-apiserver 获取 Pod/PodGroup 与节点信息;
- vc-scheduler 依据配置的调度策略为每个 Pod 选择合适的节点;
- 节点分配完成后,kubelet 从 kube-apiserver 获取 Pod 配置并启动对应容器。
JobFlow 控制器的核心实现
从 jobflow_controller.go 源码可以看到控制器的标准工作队列模式:
- 通过 Informer 监听
JobFlows(Add/Update 事件)与Jobs(Update 事件),JobTemplates 仅注册 Lister 用于读取; Run()启动 Informer Factory 并等待缓存同步,随后以wait.Until(jf.worker, time.Second, stopCh)启动 worker 循环;handleJobFlow根据当前 JobFlow 状态创建对应的状态机对象,并执行jobFlowState.Execute(req.Action);- 出错时通过限速队列
AddRateLimited重试,超过maxRequeueNum后丢弃并记录 Warning 事件。
核心同步函数syncJobFlow(jobflow_controller_action.go)按顺序完成三件事:
- 按 jobRetainPolicy 清理作业:若
JobRetainPolicy == Delete且 JobFlow 处于 Succeed 状态,删除其创建的全部 VCJob; - 按依赖顺序下发作业:调用
deployJob遍历flows,对每个 flow 检查依赖是否满足,满足则调用createJob创建 VCJob; - 汇总并更新状态:调用
getAllJobStatus收集全部 VCJob 状态,更新 JobFlow 的 status。
其中deployJob的依赖判定逻辑(jobflow_controller_action.go)为:
- flow 没有
dependsOn或targets为空:直接创建 VCJob; - 否则调用
judge检查所有目标作业是否已存在且处于Completed阶段,全部满足才创建,任一不满足则跳过(等待后续 sync 触发)。
创建的 VCJob 命名规则在 jobflow_controller_util.go 中定义:getJobName(jobFlowName, jobTemplateName)返回jobFlowName + "-" + jobTemplateName。同时,创建的 VCJob 会打上CreatedByJobFlow与CreatedByJobTemplate标签/注解,并设置 JobFlow 为 OwnerReference,这样删除 JobFlow 时会级联清理全部 VCJob。
JobTemplate:可复用的作业模板
核心语义
JobTemplate 是 VCJob 的模板,其spec直接沿用 VCJob 的 spec(JobTemplateSpec直接跟随 vcjob 的 spec)。它本身不会被 vc-controller 当作普通 VCJob 下发,而是等待被 JobFlow 引用。关键特性如下:
- JobFlow 可以引用多个 JobTemplate;
- 一个 JobTemplate 可以被多个 JobFlow 引用;
- JobTemplate 与 VCJob 可以相互转换;
- JobTemplate 简写为
jt,可通过kubectl get jt查看; - JobFlow 在引用 JobTemplate 时支持对其做 patch 修改。
JobTemplate 的增删改影响面
设计文档明确了三类操作的影响范围:
- create:创建后等待 JobFlow 使用;
- update:更新后不会影响已基于该模板创建的 VCJob,也不会影响已成功执行的 JobFlow;但可能影响尚未执行到该模板阶段的 JobFlow——尚未执行的流程会使用更新后的模板;
- delete:当 JobTemplate 正被未完成的 JobFlow 引用时,Webhook 会拦截删除请求。
JobTemplate 示例
example/jobflow/JobTemplate.yaml 给出了完整的模板定义,其 spec 与 VCJob 一致(minAvailable、schedulerName: volcano、queue、tasks等):
apiVersion: flow.volcano.sh/v1alpha1 kind: JobTemplate metadata: name: a spec: minAvailable: 1 schedulerName: volcano queue: default tasks: - replicas: 1 name: "default-nginx" template: metadata: name: web spec: containers: - image: nginx:1.14.2 command: - sh - -c - sleep 10s imagePullPolicy: IfNotPresent name: nginx resources: requests: cpu: "1" restartPolicy: OnFailureJobFlow:声明式 DAG 作业编排
字段全景
JobFlow 定义一组作业的运行流程,flows字段描述作业间的编排方式。以下为设计文档中的关键字段表(原始表格整理):
| 对象 | 属性 | 类型 | 必填 | 默认值 | 说明 |
|---|---|---|---|---|---|
| Spec | flows | Flow 数组 | 是 | — | 描述 vcjob 之间的依赖关系 |
| Spec | jobRetainPolicy | string | 是 | retain | JobFlow 成功后是否保留生成的作业(delete/retain) |
| Flow | name | string | 是 | — | 引用的 JobTemplate 名称 |
| Flow | dependsOn | DependsOn | 是 | — | JobTemplate 依赖关系 |
| Flow | patch | Patch | 否 | — | 对 JobTemplate 的 patch 修改 |
| DependsOn | targets | string 数组 | 是 | — | 当前 JobTemplate 依赖的所有 JobTemplate 名称 |
| DependsOn | probe | Probe | 否 | — | 探针类型依赖 |
| DependsOn | strategy | string | 是 | all | 依赖是否必须全部满足 |
| Probe | httpGetList | HttpGet 数组 | 否 | — | HttpGet 类型依赖 |
| Probe | tcpSocketList | TcpSocket 数组 | 否 | — | TcpSocket 类型依赖 |
| Probe | taskStatusList | TaskStatus 数组 | 否 | — | TaskStatus 类型依赖 |
| HttpGet | TaskName | string | 是 | — | vcjob 下的任务名 |
| HttpGet | Path | string | 是 | — | httpget 路径 |
| HttpGet | Port | int | 是 | — | httpget 端口 |
| HttpGet | httpHeader | HTTPHeader | 否 | — | httpget 请求头 |
| TcpSocket | TaskName | string | 是 | — | vcjob 下的任务名 |
| TcpSocket | Port | int | 是 | — | TcpSocket 端口 |
| TaskStatus | TaskName | string | 是 | — | vcjob 下的任务名 |
| TaskStatus | Phase | string | 是 | — | 任务阶段 |
Status 字段
JobFlow 的status是了解作业流运行情况的窗口,主要字段:
| 属性 | 类型 | 说明 |
|---|---|---|
pendingJobs | string 数组 | 处于 Pending 状态的 vcjob |
runningJobs | string 数组 | 处于 Running 状态的 vcjob |
failedJobs | string 数组 | 处于 Failed 状态的 vcjob |
completedJobs | string 数组 | 处于 Completed / Completing 状态的 vcjob |
terminatedJobs | string 数组 | 处于 Terminated / Terminating 状态的 vcjob |
unKnowJobs | string 数组 | 未识别状态的 vcjob |
jobStatusList | JobStatus 数组 | 所有拆分 vcjob 的状态信息(名称、状态、起止时间、重启次数、运行历史等) |
conditions | map[string]Condition | 描述所有 vcjob 的当前状态、创建时间、完成时间与信息;vcjob 状态在此额外增加 waiting 状态,用于描述依赖未满足的 vcjob |
state | State | JobFlow 的状态 |
getAllJobStatus(jobflow_controller_action.go)展示了这些字段在控制器中如何被填充:按 VCJob 的Status.State.Phase分组归入 Pending/Running/Completing/Completed/Terminating/Terminated/Failed,无法识别的进入UnKnowJobs;每个作业生成JobStatus(含RunningHistories运行历史,记录各状态的起止时间),并汇总为Conditions映射。
State 状态机
JobFlow 的状态机在 pkg/controllers/jobflow/state 中实现,factory.go的NewState根据jobFlow.Status.State.Phase分发到五种状态,实现Execute(action)接口:
| 阶段 | 状态类 | 触发语义 |
|---|---|---|
""/Pending | pendingState | JobFlow 初始状态 |
Running | runningState | 流程中存在 Running 状态的 vcjob |
Succeed | succeedState | 所有 vcjob 均达到 Completed 状态 |
Terminating | terminatingState | JobFlow 正在删除 |
Failed | failedState | 流程中存在 Failed 状态的 vcjob,后续 vcjob 无法继续下发 |
以 state/running.go 为例,Running 状态下执行SyncJobFlowAction时:
- 若
len(status.CompletedJobs) == allJobList(全部作业完成),更新为Succeed; - 若存在 Failed 或 Terminated 作业,更新为Failed。
JobFlow 状态变化遵循"作用域隔离"原则:当前 JobFlow 状态的变化不会影响其他资源。
Webhook 校验:把非法 DAG 挡在门外
JobFlow/JobTemplate 的创建与更新会经过 Admission Webhook 校验。JobFlow 的校验逻辑位于 pkg/webhooks/admission/jobflows/validate/validate_jobflow.go,其中validateJobFlowDAG将 flows 的依赖关系构造成图(graphMap),并进行两项检查:
- 同一 JobFlow 依赖中不能出现同名模板:例如
A->B->A->C中 A 出现两次,会被拒绝; - JobFlow 中不能出现闭环:例如 A → B → C → D → B 这种循环依赖,通过 DAG(有向无环图)检测拦截。
对应测试 validate_jobflow_test.go 中覆盖了"duplicate flow name"、闭环等多种非法场景。
JobTemplate 的创建校验则遵循 VCJob 参数规范(见设计文档),例如:
- job 的
minAvailable必须大于等于 0; - job 的
maxRetry必须大于等于 0; - tasks 不能为空,且不能有同名任务;
- 任务副本数不能小于 0;
- task 的
minAvailable不能大于 task 的replicas等。
此外,Webhook 还会拦截两类非预期操作:
- update jobflow:JobFlow 当前不支持更新操作,更新请求会被 Webhook 阻塞;
- delete jobflow:当 JobFlow 处于非完成状态时删除会被拦截;正常删除后,JobFlow 创建的全部 VCJob 会被直接删除(依赖 OwnerReference 级联清理)。
端到端实战:从模板到作业流
example/jobflow/README.md 给出了完整的上手步骤,前置条件是 Kubernetes 版本大于 1.17,且已安装 Volcano。
第一步:创建 JobTemplate
kubectl apply -f JobTemplate.yamlexample/jobflow/JobTemplate.yaml 中定义了 a、b、c、d、e 五个模板,每个模板内部是一个运行sleep 10s的 nginx 容器作业。
第二步:创建 JobFlow
kubectl apply -f JobFlow.yamlexample/jobflow/JobFlow.yaml 定义的依赖关系为:
apiVersion: flow.volcano.sh/v1alpha1 kind: JobFlow metadata: name: test namespace: default spec: jobRetainPolicy: delete # After jobflow runs, keep the generated job. Otherwise, delete it. flows: - name: a - name: b dependsOn: targets: ['a'] - name: c dependsOn: targets: ['b'] - name: d dependsOn: targets: ['b'] - name: e dependsOn: targets: ['c','d']这是一个典型的 DAG:
a无依赖,最先下发;b依赖a;c与d均依赖b,二者可并行运行;e依赖c与d(即strategy: all的"全部满足"语义),只有 c、d 都完成后才会启动。
第三步:查看运行状态
# 查看模板与作业流 kubectl get jt kubectl get jf # 查看作业流创建的 Pod kubectl get po由于示例中设置了jobRetainPolicy: delete,当 JobFlow 成功后,控制器会自动删除由它创建的 VCJob(对应syncJobFlow中的清理逻辑)。
使用 JobFlow 的通用流程
- 创建需要用到的 JobTemplate;
- 创建 JobFlow,其
flows字段填入用于创建 vcjob 的对应 jobtemplate; - 通过
jobRetainPolicy字段控制 JobFlow 成功后是否删除其创建的 vcjob(delete/retain,默认 retain)。
JobFlow 的 JobTemplate Patch 能力
JobFlow 在引用 JobTemplate 时支持对模板进行 patch 修改,设计文档给出的示例如下——在 flow 的patch.spec.tasks中覆盖容器命令,使该次运行的作业执行sleep 10s而不是模板默认行为:
apiVersion: flow.volcano.sh/v1alpha1 kind: JobFlow metadata: name: test namespace: default spec: jobRetainPolicy: delete flows: - name: a patch: spec: tasks: - name: "default-nginx" template: spec: containers: - name: nginx command: - sh - -c - sleep 10s功能现状与演进方向
设计文档明确列出了 JobFlow 的当前能力边界:
已实现功能:
- 创建 JobFlow 与 JobTemplate CRD;
- 支持 vcjob 顺序启动;
- 支持 vcjob 依赖其他 vcjob 启动;
- 支持 vcjob 与 JobTemplate 相互转换;
- 支持查看 JobFlow 运行状态。
尚未实现的功能(截至设计文档记录):
- JobFlow 在引用 jobtemplate 时对其修改(patch)的完善;
if语句;switch语句;for语句;- 支持 JobFlow 内的作业失败重试;
- 与 volcano-scheduler 的集成;
- 在 JobFlow 级别支持调度插件。
小结
Volcano JobFlow 用两个 CRD(JobTemplate + JobFlow)把"多作业依赖编排"沉淀为平台能力:JobTemplate 负责作业定义复用,JobFlow 以 DAG 语义驱动作业按依赖顺序自动下发。结合 jobflow_controller.go 的状态机实现、jobflow_controller_action.go 的依赖判定逻辑与 validate_jobflow.go 的 Webhook 校验,可以确认其核心机制:控制器轮询依赖目标的 Completed 状态、按 flow 顺序创建 VCJob、通过 OwnerReference 实现级联清理、以 DAG 校验保证依赖图合法性。对于 AI 训练、大数据分析等典型的多阶段批处理场景,这是一套开箱即用的声明式作业编排方案。
【免费下载链接】volcanoA Cloud Native Batch System (Project under CNCF)项目地址: https://gitcode.com/GitHub_Trending/vol/volcano
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考