Volcano JobFlow 作业编排指南:基于 DAG 依赖驱动的云原生批处理工作流引擎
2026/9/17 15:03:59 网站建设 项目流程

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 新增定义):

  1. 通过 Admission 后,kubectl在 kube-apiserver 中创建 JobTemplate 与 JobFlow(Volcano CRD)对象;
  2. JobFlowController 以 JobTemplate 为模板,根据 JobFlow 的配置与依赖规则创建对应的 VCJob;
  3. VCJob 创建后,VCJobController 根据 VCJob 配置创建对应的 Pod 与 PodGroup;
  4. Pod 与 PodGroup 创建后,vc-scheduler 从 kube-apiserver 获取 Pod/PodGroup 与节点信息;
  5. vc-scheduler 依据配置的调度策略为每个 Pod 选择合适的节点;
  6. 节点分配完成后,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)按顺序完成三件事:

  1. 按 jobRetainPolicy 清理作业:若JobRetainPolicy == Delete且 JobFlow 处于 Succeed 状态,删除其创建的全部 VCJob;
  2. 按依赖顺序下发作业:调用deployJob遍历flows,对每个 flow 检查依赖是否满足,满足则调用createJob创建 VCJob;
  3. 汇总并更新状态:调用getAllJobStatus收集全部 VCJob 状态,更新 JobFlow 的 status。

其中deployJob的依赖判定逻辑(jobflow_controller_action.go)为:

  • flow 没有dependsOntargets为空:直接创建 VCJob;
  • 否则调用judge检查所有目标作业是否已存在且处于Completed阶段,全部满足才创建,任一不满足则跳过(等待后续 sync 触发)。

创建的 VCJob 命名规则在 jobflow_controller_util.go 中定义:getJobName(jobFlowName, jobTemplateName)返回jobFlowName + "-" + jobTemplateName。同时,创建的 VCJob 会打上CreatedByJobFlowCreatedByJobTemplate标签/注解,并设置 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 一致(minAvailableschedulerName: volcanoqueuetasks等):

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: OnFailure

JobFlow:声明式 DAG 作业编排

字段全景

JobFlow 定义一组作业的运行流程,flows字段描述作业间的编排方式。以下为设计文档中的关键字段表(原始表格整理):

对象属性类型必填默认值说明
SpecflowsFlow 数组描述 vcjob 之间的依赖关系
SpecjobRetainPolicystringretainJobFlow 成功后是否保留生成的作业(delete/retain)
Flownamestring引用的 JobTemplate 名称
FlowdependsOnDependsOnJobTemplate 依赖关系
FlowpatchPatch对 JobTemplate 的 patch 修改
DependsOntargetsstring 数组当前 JobTemplate 依赖的所有 JobTemplate 名称
DependsOnprobeProbe探针类型依赖
DependsOnstrategystringall依赖是否必须全部满足
ProbehttpGetListHttpGet 数组HttpGet 类型依赖
ProbetcpSocketListTcpSocket 数组TcpSocket 类型依赖
ProbetaskStatusListTaskStatus 数组TaskStatus 类型依赖
HttpGetTaskNamestringvcjob 下的任务名
HttpGetPathstringhttpget 路径
HttpGetPortinthttpget 端口
HttpGethttpHeaderHTTPHeaderhttpget 请求头
TcpSocketTaskNamestringvcjob 下的任务名
TcpSocketPortintTcpSocket 端口
TaskStatusTaskNamestringvcjob 下的任务名
TaskStatusPhasestring任务阶段

Status 字段

JobFlow 的status是了解作业流运行情况的窗口,主要字段:

属性类型说明
pendingJobsstring 数组处于 Pending 状态的 vcjob
runningJobsstring 数组处于 Running 状态的 vcjob
failedJobsstring 数组处于 Failed 状态的 vcjob
completedJobsstring 数组处于 Completed / Completing 状态的 vcjob
terminatedJobsstring 数组处于 Terminated / Terminating 状态的 vcjob
unKnowJobsstring 数组未识别状态的 vcjob
jobStatusListJobStatus 数组所有拆分 vcjob 的状态信息(名称、状态、起止时间、重启次数、运行历史等)
conditionsmap[string]Condition描述所有 vcjob 的当前状态、创建时间、完成时间与信息;vcjob 状态在此额外增加 waiting 状态,用于描述依赖未满足的 vcjob
stateStateJobFlow 的状态

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.goNewState根据jobFlow.Status.State.Phase分发到五种状态,实现Execute(action)接口:

阶段状态类触发语义
""/PendingpendingStateJobFlow 初始状态
RunningrunningState流程中存在 Running 状态的 vcjob
SucceedsucceedState所有 vcjob 均达到 Completed 状态
TerminatingterminatingStateJobFlow 正在删除
FailedfailedState流程中存在 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),并进行两项检查:

  1. 同一 JobFlow 依赖中不能出现同名模板:例如A->B->A->C中 A 出现两次,会被拒绝;
  2. 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.yaml

example/jobflow/JobTemplate.yaml 中定义了 a、b、c、d、e 五个模板,每个模板内部是一个运行sleep 10s的 nginx 容器作业。

第二步:创建 JobFlow

kubectl apply -f JobFlow.yaml

example/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
  • cd均依赖b,二者可并行运行;
  • e依赖cd(即strategy: all的"全部满足"语义),只有 c、d 都完成后才会启动。

第三步:查看运行状态

# 查看模板与作业流 kubectl get jt kubectl get jf # 查看作业流创建的 Pod kubectl get po

由于示例中设置了jobRetainPolicy: delete,当 JobFlow 成功后,控制器会自动删除由它创建的 VCJob(对应syncJobFlow中的清理逻辑)。

使用 JobFlow 的通用流程

  1. 创建需要用到的 JobTemplate;
  2. 创建 JobFlow,其flows字段填入用于创建 vcjob 的对应 jobtemplate;
  3. 通过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),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询