用 Argo Workflow 编排 Volcano MPI 作业:在 Kubernetes 上跑通 MPI Hello World 全流程指南
2026/9/17 7:46:53 网站建设 项目流程

用 Argo Workflow 编排 Volcano MPI 作业:在 Kubernetes 上跑通 MPI Hello World 全流程指南

【免费下载链接】volcanoA Cloud Native Batch System (Project under CNCF)项目地址: https://gitcode.com/GitHub_Trending/vol/volcano

本指南完整讲解如何在 Kubernetes 集群中,以Argo Workflow 的 resource 模板(Resource Template)创建并托管Volcano Job(batch.volcano.sh/v1alpha1),从而让 Argo 的 DAG/Step 流程控制能力与 Volcano 的批量调度能力(gang scheduling、任务间依赖、ssh/svc 插件)协同工作。文中以仓库 example/integrations/argo/mpi 中的hello-world.yaml为例,带你完成从安装 Argo、准备 RBAC、提交工作流、观察日志到清理环境的完整实操,并逐段剖析 MPI master/worker 的启动时序、kubectl wait等待机制与日志跟随技巧,读完即可在你的集群中复现同样的 MPI 作业编排。

一、为什么要在 Argo Workflow 里运行 Volcano Job

Volcano 提供的是批量作业调度能力(gang scheduling、队列、任务间依赖等),而 Argo Workflow 提供的是工作流编排能力(步骤、DAG、参数化)。二者互补:Argo 负责"按什么顺序、在什么依赖关系下运行哪些作业",Volcano 负责"每个作业如何被调度、资源如何被抢占与回收"。

Argo 的resource 模板允许用户对任意 Kubernetes 资源(包括 CRD)执行 create / apply / delete / patch 等操作,因此可以天然地把 Volcano Job 作为 Argo 工作流的一个步骤来管理。仓库中 example/integrations/argo/README.md 的说明指出,通过这一机制可以将 Volcano Job 集成进 Argo Workflow,并借用 Argo 为 Volcano 增加作业依赖管理和 DAG 流程控制能力。

除了本文主角 MPI 示例,仓库还提供了另外两种编排范式的完整清单:

  • example/integrations/argo/10-job-step.yaml:使用 Argo 的Steps模板顺序编排多个 Volcano Job(step 之间串行,同一 step 内可并行多个任务)。
  • example/integrations/argo/20-job-DAG.yaml:使用 Argo 的DAG模板表达作业依赖图(如 B、C 依赖 A,D 依赖 B 与 C)。

MPI 场景之所以特殊,是因为 MPI 作业需要一个 master 节点通过 ssh 拉起多个 worker 节点协同计算,这正好同时用到了 Volcano 的 gang scheduling 与 Argo 的 DAG 并行任务(提交作业 + 跟随日志并行执行)。

二、快速开始:安装 Argo 并准备 RBAC

2.1 安装 Argo Workflows

若你尚未安装 Argo Workflows,可以使用官方 Quick Start 方式安装:

kubectl create namespace argo kubectl apply -n argo -f "https://github.com/argoproj/argo-workflows/releases/download/v3.7.10/quick-start-minimal.yaml"

该版本号与仓库示例保持一致(Argo Workflows v3.7.10);注意:以下步骤可能与 Argo Workflows 4.0+ 不兼容(见文末注意事项)。

2.2 准备 RBAC:让 Argo 能管理 Volcano Job 与 Pod

Argo 的 resource 模板在执行时相当于代表工作流绑定的 ServiceAccount 去操作目标资源,因此必须为该 ServiceAccount 授予对应资源的读写权限。仓库提供的 example/integrations/argo/mpi/rbac.yaml 定义了最小权限集(本示例基于 Argo 官方 workflow-rbac 文档改编):

apiVersion: rbac.authorization.k8s.io/v1 kind: Role metadata: name: argo-workflowtaskresults-role namespace: argo rules: - apiGroups: ["argoproj.io"] resources: ["workflowtaskresults"] verbs: ["create", "patch", "get", "list", "watch"] - apiGroups: ["batch.volcano.sh"] resources: ["jobs"] verbs: ["create", "get", "list", "watch", "delete"] - apiGroups: [""] resources: ["pods","pods/log"] verbs: ["get","list","watch"] --- apiVersion: rbac.authorization.k8s.io/v1 kind: RoleBinding metadata: name: argo-workflowtaskresults-binding namespace: argo subjects: - kind: ServiceAccount name: argo namespace: argo roleRef: kind: Role name: argo-workflowtaskresults-role apiGroup: rbac.authorization.k8s.io

要点说明:

  • batch.volcano.sh组下的jobs:Argo 需要 create / get / list / watch / delete Volcano Job(CRD),这是 resource 模板能否创建并轮询作业状态的关键;
  • podspods/log:本示例的"跟随 master 日志"任务与 master 的 InitContainer 需要通过kubectl查询 Pod 状态、读取日志,因此需要 pods 的 get / list / watch 以及 pods/log 的读取权限;
  • argoproj.io组下的workflowtaskresults:Argo 用于保存工作流任务结果,属于 Argo 自身运行所需;
  • 若你的 Volcano Job 运行在其它 namespace,请将 Role / RoleBinding 相应调整,或改用 ClusterRole / ClusterRoleBinding(对应 example/integrations/argo/README.md 中提到的"必要时可手动创建 clusterrole 与 clusterrolebinding")。

安装:

kubectl apply -f rbac.yaml

三、提交 MPI Hello World 工作流

3.1 创建 WorkflowTemplate 并提交

仓库的 example/integrations/argo/mpi/hello-world.yaml 定义的是一个名为volcano-mpi-helloWorkflowTemplate(而不是直接实例化的 Workflow)。创建模板后,使用argo submit --from workflowtemplate/...实例化并提交:

kubectl apply -f hello-world.yaml WF_NAME=volcano-mpi-hello-$(date +%s) argo submit --from workflowtemplate/volcano-mpi-hello -n argo --name "$WF_NAME" argo watch "$WF_NAME" -n argo

argo watch会实时跟踪工作流直至结束,终端中即可看到两个 DAG 任务(submit-jobfollow-logs)的状态流转。

3.2 参数化说明

WorkflowTemplate 在spec.arguments.parameters中声明了一个参数:

arguments: parameters: - name: job-name value: mpi-hello-job

job-name同时被用于:

  • Volcano Job 的metadata.name
  • 后续kubectl查询 Pod 时使用的标签选择器volcano.sh/job-name={{workflow.parameters.job-name}}

这意味着你可以在提交时用-p job-name=my-custom-name覆盖作业名,多个 MPI 作业即可在同一个 namespace 内并行运行而互不干扰。

四、剖析 hello-world.yaml:WorkflowTemplate 内部结构

hello-world.yaml的整体结构是一个包含 3 个模板的 WorkflowTemplate:入口main(DAG)、follow-master-logs(日志跟随容器)与submit-volcano-job(resource 模板创建 Volcano Job)。下面逐层拆解。

4.1 入口模板:DAG 并行编排

spec: serviceAccountName: argo entrypoint: main templates: - name: main dag: tasks: - name: submit-job template: submit-volcano-job - name: follow-logs template: follow-master-logs

入口模板是一个 DAG,其中包含两个互相之间没有依赖的任务:

  • submit-job:调用 resource 模板创建 Volcano Job,并等待其状态变为Completed/Failed
  • follow-logs:与submit-job并行运行,负责在 MPI master Pod 就绪后持续流式输出其日志。

这正是本示例的核心设计之一:submit-job负责"等作业跑完",follow-logs负责"实时展示日志",二者并行使得用户既能拿到最终结果,又能实时看到 MPI 进程的输出(否则资源模板等待期间日志不可见)。

4.2 resource 模板:以 Argo 管理 Volcano Job 生命周期

- name: submit-volcano-job resource: action: create setOwnerReference: true successCondition: status.state.phase == Completed failureCondition: status.state.phase == Failed manifest: | apiVersion: batch.volcano.sh/v1alpha1 kind: Job ...

关键字段:

  • action: create:对下面manifest中的资源执行 create(也支持 apply / delete / patch);
  • setOwnerReference: true:自动为创建的 Volcano Job 设置 owner reference 指向当前 Workflow,从而让作业随工作流生命周期被清理(对应 example/integrations/argo/README.md 中"为确保 Argo 能管理其创建的资源需添加 ownerReferences"的说明;若手动声明,写法为ownerReferences: [{apiVersion: argoproj.io/v1alpha1, blockOwnerDeletion: true, kind: Workflow, name: "{{workflow.name}}", uid: "{{workflow.uid}}"}]);
  • successCondition/failureCondition:Argo 会通过kubectl get -o json -w轮询资源状态,并按 Kubernetes 标签选择语法对任意字段求值。此处对应 Volcano Job 状态中的status.state.phase
    • == Completed:作业所有任务均已完成(成功);
    • == Failed:作业重试达到上限后失败;
    • 两者可根据实际业务场景调整。

4.3 Volcano Job 主体:MPI master / worker 拓扑

被 Argo 托管的 Volcano Job 本体(batch.volcano.sh/v1alpha1Job)如下:

spec: minAvailable: 3 schedulerName: volcano plugins: ssh: [] svc: [] tasks: - name: mpimaster replicas: 1 policies: - event: TaskCompleted action: CompleteJob template: spec: serviceAccountName: argo initContainers: - name: wait-for-workers image: mcr.microsoft.com/oss/kubernetes/kubectl:v1.26.3 command: ["/bin/bash", "-c", "kubectl wait pod -l volcano.sh/job-name={{workflow.parameters.job-name}},volcano.sh/task-spec=mpiworker --for=condition=Ready --timeout=600s"] containers: - name: mpimaster image: volcanosh/example-mpi:0.0.3 ... - name: mpiworker replicas: 2 template: spec: containers: - name: mpiworker image: volcanosh/example-mpi:0.0.3 ...

逐一解读:

  • minAvailable: 3:1 个 mpimaster + 2 个 mpiworker。这是 gang scheduling 的"最小可用成员数",只有 3 个任务都有资源可调度时才会被整体拉起(参见 docs/min-avaiable-member-resource.md 与 docs/gang-aware-eviction-design.md 的相关设计说明);
  • schedulerName: volcano:明确让 Volcano 调度器接管该作业的 Pod 调度;
  • plugins: {ssh: [], svc: []}:启用 Volcano 的ssh 插件(在/etc/volcano下生成各任务间 ssh 所需的 host 文件,如mpiworker.host)与svc 插件(为每个任务创建 headless Service,使 Pod 能以<pod>.<task>-<index>.<job>形式互相发现,对应日志中mpi-hello-job-mpiworker-0.mpi-hello-job这样的 DNS 名)。原生 MPI 示例见 example/integrations/mpi/mpi-example.yaml,其 master 命令与 hello-world.yaml 一脉相承:
    MPI_HOST=`cat /etc/volcano/mpiworker.host | tr "\n" ","`; mkdir -p /var/run/sshd; /usr/sbin/sshd; mpiexec --allow-run-as-root --host ${MPI_HOST} -np 2 mpi_hello_world;
  • master 的policiesevent: TaskCompleted+action: CompleteJob——master 任务完成后即宣告整个 Job 完成(status.state.phase = Completed),这正是 resource 模板 successCondition 得以满足的信号;
  • master 的 InitContainerwait-for-workers:由于 master 要通过 ssh 连接 worker,必须先等 2 个 worker Pod 的 22 端口就绪(readinessProbe为 TCP 探测 22 端口)。InitContainer 使用kubectl wait pod -l volcano.sh/job-name=...,volcano.sh/task-spec=mpiworker --for=condition=Ready --timeout=600s等待;这要求 InitContainer 具备查询 Pod 的 RBAC 权限(即前文 rbac.yaml 中 pods get/list/watch 的来源);
  • resources.requests/limitsnvidia.com/gpu: 0rdma/ib: 0:预留的扩展资源位。仓库注释说明:可将这些值从 0 调大以消耗节点上的 NIC 或 GPU(如配合 SR-IOV device plugin),从而把 MPI 计算固定到特定节点;当一个节点上的全部 NIC / GPU 都被本作业消耗后,其它 Pod 将无法再调度到该节点,适用于独占 RDMA / GPU 卡的 MPI 集群场景;
  • worker 的readinessProbe:TCP 探测 22 端口,initialDelaySeconds: 5periodSeconds: 10failureThreshold: 5,用于向kubectl wait --for=condition=Ready暴露就绪状态。

4.4 可选的 dependsOn 编排

master 任务的模板中有一段被注释掉的配置:

# Optional: use dependsOn to start master after workers exist. # If dependsOn is enabled in this example, set minAvailable to 2 # (worker replica count). Also bear in mind that master is then # not gang-scheduled together with workers. # dependsOn: # name: # - mpiworker

若启用dependsOn: mpiworker,master 将在 worker 之后才被创建,此时:

  • minAvailable应改为 2(即只对 worker 做 gang 调度);
  • master 不再与 worker 一起被 gang 调度,失去了整体性保证。

默认(不启用)时,master 与 worker 作为一个整体被 gang 调度,通过 InitContainer +kubectl wait实现"先等 worker 就绪、再启动 ssh 计算"的时序,这是本示例推荐的做法。

五、日志跟随机制与预期输出

5.1 follow-master-logs 任务的实现

- name: follow-master-logs serviceAccountName: argo activeDeadlineSeconds: 300 container: image: mcr.microsoft.com/oss/kubernetes/kubectl:v1.26.3 command: [bash, -c] args: - | NS=$(cat /var/run/secrets/kubernetes.io/serviceaccount/namespace) echo "Waiting for MPI master pod..." POD="" until [ -n "$POD" ]; do POD=$(kubectl get pods -n $NS \ -l volcano.sh/job-name={{workflow.parameters.job-name}},volcano.sh/task-spec=mpimaster \ -o jsonpath='{.items[0].metadata.name}' 2>/dev/null) sleep 2 done echo "Found pod $POD" while true; do PHASE=$(kubectl get pod $POD -n $NS -o jsonpath='{.status.phase}' 2>/dev/null) if [[ "$PHASE" == "Running" || "$PHASE" == "Succeeded" || "$PHASE" == "Failed" ]]; then break; fi sleep 2 done echo "Streaming logs..." kubectl logs -f -n $NS -c mpimaster $POD || true

该任务共分三步:

  1. 轮询查找 master Pod:通过标签选择器volcano.sh/job-name=<作业名>+volcano.sh/task-spec=mpimaster等待 master Pod 出现(每 2 秒一次);
  2. 等待 Pod 进入终态或运行态:直到 phase 为 Running / Succeeded / Failed;
  3. 流式输出日志kubectl logs -f跟随 master 容器(容器名mpimaster)输出,|| true保证日志流中断不至于令任务失败。

activeDeadlineSeconds: 300为其设置了 5 分钟的执行上限。这套"kubectl wait + 轮询 + 日志跟随"的思路来源于 Azure 面向 Kubernetes 上 MPI/HPC 的实践(NDM v4 A100 部署指南),仓库将其引入 Argo 模板。

5.2 预期输出

follow-master-logs任务的预期输出如下:

Waiting for MPI master pod... Found pod mpi-hello-job-mpimaster-0 Waiting for pod to start or finish... Pod phase: Running Streaming logs... Running MPI hello world... Warning: Permanently added 'mpi-hello-job-mpiworker-0.mpi-hello-job' (ED25519) to the list of known hosts. Warning: Permanently added 'mpi-hello-job-mpiworker-1.mpi-hello-job' (ED25519) to the list of known hosts. Hello world from processor mpi-hello-job-mpiworker-0, rank 0 out of 2 processors Hello world from processor mpi-hello-job-mpiworker-1, rank 1 out of 2 processors

可以看到:

  • 日志中出现两个 worker 的 DNS 名mpi-hello-job-mpiworker-0.mpi-hello-job/mpi-hello-job-mpiworker-1.mpi-hello-job,印证了 svc 插件生成的 headless Service 命名规则;
  • 两条 "Hello world" 表明mpiexec -np 2确实通过 ssh 在两个 worker 上启动了 2 个 MPI rank;
  • 当 master 完成输出后,TaskCompleted策略将 Job 置为 Completed,submit-job的 successCondition 随即满足,整个 Workflow 以成功收尾。

六、清理环境

运行完成后按以下顺序清理资源:

argo delete "$WF_NAME" -n argo kubectl delete -f hello-world.yaml kubectl delete -f rbac.yaml

由于setOwnerReference: true的存在,Volcano Job 会随工作流删除而被级联清理;argo delete删除工作流实例,随后删除 WorkflowTemplate 与 RBAC。

七、进阶:用 Steps 与 DAG 编排多个 Volcano Job

理解 MPI 单作业流程后,可进一步查看仓库中另外两个完整示例,将 Volcano Job 接入更复杂的流水线:

  • Steps 顺序编排:example/integrations/argo/10-job-step.yaml 定义了volcano-step-job流程:第一步hello-1执行完毕后,第二步并行运行hello-2ahello-2b,每个步骤都是一个独立的 Volcano Job(generateName: step-job-<task>-),通过ownerReferences绑定工作流生命周期;
  • DAG 依赖编排:example/integrations/argo/20-job-DAG.yaml 定义了volcano-dag-job:任务 B、C 依赖 A,任务 D 依赖 B 与 C,即"先 A 后 B/C,最后 D"的分阶段批量作业依赖图。

两个示例的 Volcano Job 均带有:

policies: - event: PodEvicted action: RestartJob maxRetry: 1 queue: default plugins: ssh: [] env: [] svc: []

这些字段展示了 Volcano Job 面向长时/容错场景的常用配置:Pod 被驱逐时重启作业、失败重试上限、使用default队列,以及同时启用 ssh / env / svc 插件。若作业为长期运行型(如 nginx),example/integrations/argo/README.md 特别提醒:必须为模板设置activeDeadlineSeconds,否则 Workflow 将因资源模板一直等待而无法进入下一步。

八、注意事项与适用前提

根据仓库 example/integrations/argo/mpi/README.md 的 Notes,使用本方案时有以下几点限制:

  1. Argo 版本兼容性:上述步骤基于 Argo Workflows v3.7.10 编写,可能与Argo Workflows 4.0+不兼容(resource 模板的字段语义可能变化,请以实际版本文档为准);
  2. RBAC 要求:InitContainer 使用kubectl wait查询 worker Pod 状态,因此必须为其授予额外的 Pod 查询权限(即rbac.yaml中 pods get/list/watch),否则等待会失败;
  3. 资源独占:将hello-world.yaml中的nvidia.com/gpurdma/ib请求从 0 调大可消耗节点上的 NIC / GPU;当节点资源被本作业完全占用后,其它 Pod 将无法调度到该节点,请谨慎规划资源配额;
  4. 适用前提:本示例需要集群中已安装 Volcano 调度器(schedulerName: volcano依赖其工作),且 Argo 与示例资源位于argonamespace(WorkflowTemplate.metadata.namespace、RoleBinding 均指向argo),若部署到其它 namespace 需同步调整。

九、进一步阅读

  • 原生(非 Argo)Volcano MPI 作业示例:example/integrations/mpi/mpi-example.yaml(hello-world.yaml 中 master/worker 语法即改编自此处);
  • Argo 集成 Volcano Job 总览(RBAC 与 ownerReferences 详解):example/integrations/argo/README.md;
  • Steps / DAG 编排完整清单:example/integrations/argo/10-job-step.yaml、example/integrations/argo/20-job-DAG.yaml;
  • Volcano 作业 API 与任务依赖设计:docs/job-api.md、docs/task-order.md、docs/gang-aware-eviction-design.md;
  • 最小可用成员资源(minAvailable)设计:docs/min-avaiable-member-resource.md。

通过本指南,你已经掌握了一条完整可运行的"Argo 提交 → Volcano 调度 → MPI 计算 → 日志回流"链路:既能用 Argo 的 DAG/Steps 表达批量作业的依赖与并发,又能享受 Volcano 对 MPI 这类多角色作业的 gang 调度、ssh/svc 插件与任务策略带来的开箱即用能力。

【免费下载链接】volcanoA Cloud Native Batch System (Project under CNCF)项目地址: https://gitcode.com/GitHub_Trending/vol/volcano

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

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

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

立即咨询