分片有三个,Pod 却只起了一个
目标
用 Argo Workflows 真实运行一条管道,其中前面步骤产生的数据决定后面步骤的形态。通过节点记录来确认:用脚本输出做扇出(fan-out)、按条件分支、 汇总分散的结果、吸收失败、等待审批,以及在工作流之间加锁。
为什么重要
数据处理管道在运行之前,并不知道要处理多少个分片。必须先运行拆分输入的步骤才能确定分片数, 然后生成相应数量的并行任务。在 YAML 中预先写死列表的方式无法表达这一点,所以 Argo Workflows 提供了: 把一个节点的输出用作下一个节点循环列表的 withParam、对每个元素分别求值的 when,以及扇出输出的自动聚合。 作为代价,并行是默认行为,所以对使用共享资源的步骤,要自己设置上限和锁,一个分片的失败是否算作整体失败,也要在设计中决定。 考试会问这些模板类型和 spec 字段在执行时会改变什么。
步骤
- 在
/root/capa-data/split.yaml中编写 Workflow,提交到 argo 命名空间并等待其结束。generateName 为split-,标签为capa-data: split,entrypoint DAG 为main,其中 script 模板split(镜像alpine:3.20,command[sh])在标准输出中只输出一行 JSON 列表[{"id":"a","size":3},{"id":"b","size":12},{"id":"c","size":7}],之后任务peek把该结果作为参数shards接收,并用 busybox:1.36 输出。 - 在
/root/capa-data/fanout.yaml中编写并提交 Workflow(generateNamefanout-,标签capa-data: fanout)。在 entrypoint DAGmain中,任务split(第 1 步的 script 模板)之后,任务process要用 withParam 接收split的结果列表,对每个元素以参数id、size运行容器模板process(busybox:1.36)。工作流必须为 Succeeded。 - 在
/root/capa-data/serial.yaml中创建与第 2 步相同的扇出 Workflow,generateName 为serial-,标签为capa-data: serial,但让process容器执行sleep 3,并在 DAG 模板main中设置parallelism: 1后提交。三个processPod 的运行区间不得相互重叠。 - 在
/root/capa-data/branch.yaml中编写并提交 Workflow(generateNamebranch-,标签capa-data: branch)。split之后的两个任务big、small都用 withParam 接收同一个列表,big通过when只对 size 大于 10 的元素、small只对 size 不超过 10 的元素,运行容器模板work(参数id、lane)。条件为假的元素对应的节点必须保留为 Skipped。 - 在
/root/capa-data/agg.yaml中编写并提交 Workflow(generateNameagg-,标签capa-data: agg)。扇出任务count对每个元素,用文件统计seq 1 <size>的行数,并作为输出参数lines(valueFrom.path)输出;任务total把{{tasks.count.outputs.parameters.lines}}作为参数values接收,用 script 计算总和写入文件,并作为输出参数sum输出。total节点的输出参数 sum 必须等于三个 size 之和。 - 在
/root/capa-data/tolerate.yaml中编写并提交 Workflow(generateNametolerate-,标签capa-data: tolerate)。扇出任务check运行一个容器,size 超过 10 时以退出码 3 失败,并用continueOn吸收该失败,之后运行任务report。分片 b 的节点必须为 Failed,工作流和 report 必须为 Succeeded。 - 在
/root/capa-data/approve.yaml中,编写按split→ suspend 模板approve→ 容器模板publish顺序的 steps Workflow(generateNameapprove-,标签capa-data: approve),并在不带--wait的情况下提交。确认工作流停在 approve,等待 10 秒以上后用argo resume恢复,使其以 Succeeded 结束。 - 在 argo 命名空间中创建 ConfigMap
capa-data-locks(键warehouse的值为"1"),并在/root/capa-data/locked.yaml中编写用 spec.synchronization 把该键用作信号量的 Workflow(generateNamelocked-,标签capa-data: locked,容器执行sleep 8)。用同一个文件连续提交两个工作流,让两者都以 Succeeded 结束。两个工作流的运行区间不得重叠。
参考
- VM 中有 k3s、Argo Workflows v4.1.3(argo 命名空间)和 argo CLI。工作流以 argo 命名空间的 default 账户运行。
busybox:1.36、alpine:3.20已预先下载。 - 提交并等待:
argo submit -n argo <파일> --wait(占位符为文件名),查看节点:argo get -n argo <이름>(占位符为名称),列表:argo list -n argo -l capa-data=<값>(占位符为标签值)。 - 评分器会查看每个步骤标签下最新的工作流。修改后重新提交,就会以新的工作流来评分。
- 常见错误:在 withParam 中混用
{{steps...}}和{{tasks...}}。DAG 里用 tasks,steps 里用 steps。 - 常见错误:when 表达式中的字符串不加引号就比较。数字比较保持原样,字符串比较要用单引号括起来。
- Loops(withParam)、Conditionals、Suspending、Synchronization、Field Reference(continueOn)
脚本的标准输出变成数据
在 /root/capa-data/split.yaml 中编写 Workflow,提交到 argo 命名空间并等待其结束。generateName 为 split-,标签为 capa-data: split,entrypoint DAG 为 main,其中 script 模板 split(镜像 alpine:3.20,command [sh])在标准输出中只输出一行 JSON 列表 [{"id":"a","size":3},{"id":"b","size":12},{"id":"c","size":7}],之后任务 peek 把该结果作为参数 shards 接收,并用 busybox:1.36 输出。
script 模板会把 source 写成文件并用 command 执行,并把标准输出作为节点的 outputs.result 输出。后面的任务用 {{tasks.<이름>.outputs.result}}(占位符为任务名称)读取。在此版本(v4.1.3)中,没有任何地方引用的 script 的 result 不会留在节点记录里(实测)。用 argo submit --wait 提交会一直等到结束。
前一步列表有多长,就生成多少个 Pod
在 /root/capa-data/fanout.yaml 中编写并提交 Workflow(generateName fanout-,标签 capa-data: fanout)。在 entrypoint DAG main 中,任务 split(第 1 步的 script 模板)之后,任务 process 要用 withParam 接收 split 的结果列表,对每个元素以参数 id、size 运行容器模板 process(busybox:1.36)。工作流必须为 Succeeded。
DAG 任务的结果用 {{tasks.<이름>.outputs.result}}(占位符为任务名称)读取。withParam 元素的字段用 {{item.<필드>}}(占位符为字段名称)读取。withItems 是把列表预先写在 YAML 中的方式,所以列表长度不会在运行中确定。
让扇出一次只放行一个
在 /root/capa-data/serial.yaml 中创建与第 2 步相同的扇出 Workflow,generateName 为 serial-,标签为 capa-data: serial,但让 process 容器执行 sleep 3,并在 DAG 模板 main 中设置 parallelism: 1 后提交。三个 process Pod 的运行区间不得相互重叠。
parallelism 既可以设在工作流 spec 整体上,也可以设在单个模板上。评分器用节点的 startedAt、finishedAt 来检查是否重叠。
按大小分送到不同的分支
在 /root/capa-data/branch.yaml 中编写并提交 Workflow(generateName branch-,标签 capa-data: branch)。split 之后的两个任务 big、small 都用 withParam 接收同一个列表,big 通过 when 只对 size 大于 10 的元素、small 只对 size 不超过 10 的元素,运行容器模板 work(参数 id、lane)。条件为假的元素对应的节点必须保留为 Skipped。
when 会把参数替换完成后的字符串当作表达式求值。在一个任务中同时使用 withParam 和 when,会对每个元素分别求值。
把分散的结果重新汇总
在 /root/capa-data/agg.yaml 中编写并提交 Workflow(generateName agg-,标签 capa-data: agg)。扇出任务 count 对每个元素,用文件统计 seq 1 <size> 的行数,并作为输出参数 lines(valueFrom.path)输出;任务 total 把 {{tasks.count.outputs.parameters.lines}} 作为参数 values 接收,用 script 计算总和写入文件,并作为输出参数 sum 输出。total 节点的输出参数 sum 必须等于三个 size 之和。
在后面的任务中读取扇出任务的输出参数时,各元素的值会汇集成 JSON 列表字符串。从该字符串中只取出数字再相加即可。评分器除了总和,还会查看 total 收到的列表。
吸收一个分片的失败,并继续到报告
在 /root/capa-data/tolerate.yaml 中编写并提交 Workflow(generateName tolerate-,标签 capa-data: tolerate)。扇出任务 check 运行一个容器,size 超过 10 时以退出码 3 失败,并用 continueOn 吸收该失败,之后运行任务 report。分片 b 的节点必须为 Failed,工作流和 report 必须为 Succeeded。
retryStrategy 是重新尝试同一件事,continueOn 是承认失败后继续往下走。请比较扇出中的一个元素失败时,DAG 的依赖任务会怎样。
停住,直到有人批准
在 /root/capa-data/approve.yaml 中,编写按 split → suspend 模板 approve → 容器模板 publish 顺序的 steps Workflow(generateName approve-,标签 capa-data: approve),并在不带 --wait 的情况下提交。确认工作流停在 approve,等待 10 秒以上后用 argo resume 恢复,使其以 Succeeded 结束。
不给 suspend 模板设置 duration,就会一直等到有人恢复。暂停期间,请确认工作流的 phase 和 approve 节点的 phase。评分器会查看 approve 节点停留的时间。
不让两个工作流同时写入同一个仓储
在 argo 命名空间中创建 ConfigMap capa-data-locks(键 warehouse 的值为 "1"),并在 /root/capa-data/locked.yaml 中编写用 spec.synchronization 把该键用作信号量的 Workflow(generateName locked-,标签 capa-data: locked,容器执行 sleep 8)。用同一个文件连续提交两个工作流,让两者都以 Succeeded 结束。两个工作流的运行区间不得重叠。
parallelism 是一个工作流内部的上限,信号量则是在工作流之间共享的锁。第二个工作流等待期间,请查看 status.synchronization 和 message。此版本(v4.1.3)会把旧的单数字段 semaphore 当作未知字段拒绝,只接受列表字段(semaphores、mutexes)(实测)。