ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

Volcano JobFlow 作业编排指南:基于 DAG 依赖驱动的云原生批处理工作流引擎

Volcano JobFlow 作业编排指南:基于 DAG 依赖驱动的云原生批处理工作流引擎 Volcano JobFlow 作业编排指南基于 DAG 依赖驱动的云原生批处理工作流引擎【免费下载链接】volcanoA Cloud Native Batch System (Project under CNCF)项目地址: https://gitcode.com/GitHub_Trending/vol/volcanoVolcano 作为 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.mdJobTemplate作业模板缩写 jtVCJob 的模板定义作业的完整 spec但不被 job 控制器直接下发等待被 JobFlow 引用JobFlow作业流缩写 jf定义一组作业的运行流程通过flows字段描述作业之间的依赖关系顺序执行、并行执行、条件依赖等。JobFlow 不是通用工作流引擎它了解 VCJob 的细节因此能够为用户提供远超通用引擎的作业感知能力例如作业运行状态、起止时间戳、下一个要运行的作业、Pod 失败比率等。从设计文档的 Scope 划分来看In ScopeJobFlow 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 与 JobFlowVolcano CRD对象JobFlowController 以 JobTemplate 为模板根据 JobFlow 的配置与依赖规则创建对应的 VCJobVCJob 创建后VCJobController 根据 VCJob 配置创建对应的 Pod 与 PodGroupPod 与 PodGroup 创建后vc-scheduler 从 kube-apiserver 获取 Pod/PodGroup 与节点信息vc-scheduler 依据配置的调度策略为每个 Pod 选择合适的节点节点分配完成后kubelet 从 kube-apiserver 获取 Pod 配置并启动对应容器。JobFlow 控制器的核心实现从 jobflow_controller.go 源码可以看到控制器的标准工作队列模式通过 Informer 监听JobFlowsAdd/Update 事件与JobsUpdate 事件JobTemplates 仅注册 Lister 用于读取Run()启动 Informer Factory 并等待缓存同步随后以wait.Until(jf.worker, time.Second, stopCh)启动 worker 循环handleJobFlow根据当前 JobFlow 状态创建对应的状态机对象并执行jobFlowState.Execute(req.Action)出错时通过限速队列AddRateLimited重试超过maxRequeueNum后丢弃并记录 Warning 事件。核心同步函数syncJobFlowjobflow_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 的 specJobTemplateSpec直接跟随 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字段描述作业间的编排方式。以下为设计文档中的关键字段表原始表格整理对象属性类型必填默认值说明SpecflowsFlow 数组是—描述 vcjob 之间的依赖关系SpecjobRetainPolicystring是retainJobFlow 成功后是否保留生成的作业delete/retainFlownamestring是—引用的 JobTemplate 名称FlowdependsOnDependsOn是—JobTemplate 依赖关系FlowpatchPatch否—对 JobTemplate 的 patch 修改DependsOntargetsstring 数组是—当前 JobTemplate 依赖的所有 JobTemplate 名称DependsOnprobeProbe否—探针类型依赖DependsOnstrategystring是all依赖是否必须全部满足ProbehttpGetListHttpGet 数组否—HttpGet 类型依赖ProbetcpSocketListTcpSocket 数组否—TcpSocket 类型依赖ProbetaskStatusListTaskStatus 数组否—TaskStatus 类型依赖HttpGetTaskNamestring是—vcjob 下的任务名HttpGetPathstring是—httpget 路径HttpGetPortint是—httpget 端口HttpGethttpHeaderHTTPHeader否—httpget 请求头TcpSocketTaskNamestring是—vcjob 下的任务名TcpSocketPortint是—TcpSocket 端口TaskStatusTaskNamestring是—vcjob 下的任务名TaskStatusPhasestring是—任务阶段Status 字段JobFlow 的status是了解作业流运行情况的窗口主要字段属性类型说明pendingJobsstring 数组处于 Pending 状态的 vcjobrunningJobsstring 数组处于 Running 状态的 vcjobfailedJobsstring 数组处于 Failed 状态的 vcjobcompletedJobsstring 数组处于 Completed / Completing 状态的 vcjobterminatedJobsstring 数组处于 Terminated / Terminating 状态的 vcjobunKnowJobsstring 数组未识别状态的 vcjobjobStatusListJobStatus 数组所有拆分 vcjob 的状态信息名称、状态、起止时间、重启次数、运行历史等conditionsmap[string]Condition描述所有 vcjob 的当前状态、创建时间、完成时间与信息vcjob 状态在此额外增加 waiting 状态用于描述依赖未满足的 vcjobstateStateJobFlow 的状态getAllJobStatusjobflow_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)接口阶段状态类触发语义/PendingpendingStateJobFlow 初始状态RunningrunningState流程中存在 Running 状态的 vcjobSucceedsucceedState所有 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并进行两项检查同一 JobFlow 依赖中不能出现同名模板例如A-B-A-C中 A 出现两次会被拒绝JobFlow 中不能出现闭环例如 A → B → C → D → B 这种循环依赖通过 DAG有向无环图检测拦截。对应测试 validate_jobflow_test.go 中覆盖了duplicate flow name、闭环等多种非法场景。JobTemplate 的创建校验则遵循 VCJob 参数规范见设计文档例如job 的minAvailable必须大于等于 0job 的maxRetry必须大于等于 0tasks 不能为空且不能有同名任务任务副本数不能小于 0task 的minAvailable不能大于 task 的replicas等。此外Webhook 还会拦截两类非预期操作update jobflowJobFlow 当前不支持更新操作更新请求会被 Webhook 阻塞delete jobflow当 JobFlow 处于非完成状态时删除会被拦截正常删除后JobFlow 创建的全部 VCJob 会被直接删除依赖 OwnerReference 级联清理。端到端实战从模板到作业流example/jobflow/README.md 给出了完整的上手步骤前置条件是 Kubernetes 版本大于 1.17且已安装 Volcano。第一步创建 JobTemplatekubectl apply -f JobTemplate.yamlexample/jobflow/JobTemplate.yaml 中定义了 a、b、c、d、e 五个模板每个模板内部是一个运行sleep 10s的 nginx 容器作业。第二步创建 JobFlowkubectl 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]这是一个典型的 DAGa无依赖最先下发b依赖ac与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 成功后是否删除其创建的 vcjobdelete/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 用两个 CRDJobTemplate 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),仅供参考
返回列表