シャードは 3 つなのに Pod は 1 つしか起動しなかった
目標
前のステップが作ったデータで、後ろのステップの形が決まるパイプラインを、Argo Workflowsで実際に実行します。スクリプトの出力でファンアウトし、条件でブランチを分け、散らばった結果を集め、失敗を吸収し、承認を待ち、ワークフローの間にロックをかけることを、ノードの記録で確認します。
なぜ重要なのか
データ処理のパイプラインは、実行する前には、いくつのチャンクを処理するのかがわかりません。入力を分割するステップが動いてはじめて、チャンクの数が決まり、その数だけ並列の作業ができる必要があります。YAMLにリストをあらかじめ書く方式では、これを表現できないため、Argo Workflowsは、あるノードの出力を、次のノードの繰り返しリストとして使うwithParam、要素ごとに評価されるwhen、ファンアウトの出力の自動集計を用意しています。その代わり、並列性が基本なので、共有リソースを使うステップには、上限とロックを自分でかける必要があり、チャンク1つの失敗を、全体の失敗と見なすかどうかも、設計で決める必要があります。試験は、このテンプレートの種類とspecのフィールドが、実行で何を変えるかを問います。
ステップ
/root/capa-data/split.yamlにWorkflowを作成して、argoネームスペースに提出し、終わるまで待ってください。generateNameはsplit-、ラベルはcapa-data: split、entrypointのDAGmainで、scriptテンプレートsplit(イメージalpine:3.20、command[sh])が、標準出力にJSONのリスト[{"id":"a","size":3},{"id":"b","size":12},{"id":"c","size":7}]を1行だけ出力し、そのあとのタスクpeekが、その結果をパラメーターshardsとして受け取って、busybox:1.36で出力します。/root/capa-data/fanout.yamlに、Workflow(generateNamefanout-、ラベルcapa-data: fanout)を作成して、提出してください。entrypointのDAGmainで、タスクsplit(ステップ1のscriptテンプレート)のあとに、タスクprocessが、splitの結果のリストをwithParamで受け取り、要素ごとに、パラメーターid・sizeで、コンテナテンプレートprocess(busybox:1.36)を実行する必要があります。ワークフローがSucceededである必要があります。/root/capa-data/serial.yamlに、ステップ2と同じファンアウトのWorkflowを、generateNameserial-、ラベルcapa-data: serialで作成し、processコンテナがsleep 3を実行し、DAGテンプレートmainにparallelism: 1を置いて、提出してください。3つのprocessのPodの実行区間が、互いに重ならない必要があります。/root/capa-data/branch.yamlに、Workflow(generateNamebranch-、ラベルcapa-data: branch)を作成して、提出してください。splitのあとに、2つのタスクbig・smallが、どちらも同じリストをwithParamで受け取り、bigはwhenでsizeが10を超える要素だけを、smallは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が、3つの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に、stepsで、split→ suspendテンプレートapprove→ コンテナテンプレートpublishの順序の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を実行)を作成してください。同じファイルで、ワークフロー2つを続けて提出して、どちらもSucceededで終わらせます。2つのワークフローの実行区間は、重ならない必要があります。
参考
- 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のDAGmainで、scriptテンプレートsplit(イメージalpine:3.20、command[sh])が、標準出力にJSONのリスト[{"id":"a","size":3},{"id":"b","size":12},{"id":"c","size":7}]を1行だけ出力し、そのあとのタスク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(generateNamefanout-、ラベルcapa-data: fanout)を作成して、提出してください。entrypointのDAGmainで、タスクsplit(ステップ1のscriptテンプレート)のあとに、タスクprocessが、splitの結果のリストをwithParamで受け取り、要素ごとに、パラメーターid・sizeで、コンテナテンプレートprocess(busybox:1.36)を実行する必要があります。ワークフローがSucceededである必要があります。
DAGのタスクの結果は、{{tasks.<이름>.outputs.result}}(プレースホルダーはタスク名です)で読み取ります。withParamの要素のフィールドは、{{item.<필드>}}(プレースホルダーはフィールド名です)です。withItemsは、リストをYAMLにあらかじめ書いておく方式なので、リストの長さが、実行中には決まりません。
ファンアウトを1回に1つずつ流す
/root/capa-data/serial.yamlに、ステップ2と同じファンアウトのWorkflowを、generateNameserial-、ラベルcapa-data: serialで作成し、processコンテナがsleep 3を実行し、DAGテンプレートmainにparallelism: 1を置いて、提出してください。3つのprocessのPodの実行区間が、互いに重ならない必要があります。
parallelismは、ワークフローのspec全体にも、テンプレート1つにも置けます。採点は、ノードのstartedAt・finishedAtで、重なりを見ます。
サイズに応じて別のブランチへ送る
/root/capa-data/branch.yamlに、Workflow(generateNamebranch-、ラベルcapa-data: branch)を作成して、提出してください。splitのあとに、2つのタスクbig・smallが、どちらも同じリストをwithParamで受け取り、bigはwhenでsizeが10を超える要素だけを、smallは10以下の要素だけを、コンテナテンプレートwork(パラメーターid・lane)で実行します。条件が偽の要素のノードは、Skippedとして残る必要があります。
whenは、パラメーターの置換が終わった文字列を、式として評価します。withParamとwhenを1つのタスクに一緒に置くと、要素ごとに別々に評価されます。
散らばった結果を、もう一度集めて合算する
/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が、3つのsizeの合計である必要があります。
ファンアウトのタスクの出力パラメーターを、後ろのタスクで読むと、要素ごとの値が、JSONのリストの文字列として集まります。その文字列から数字だけを取り出して足せばよいです。採点は、合計とあわせて、totalが受け取ったリストも見ます。
1つのチャンクの失敗を吸収して、レポートまで進む
/root/capa-data/tolerate.yamlに、Workflow(generateNametolerate-、ラベルcapa-data: tolerate)を作成して、提出してください。ファンアウトのタスクcheckは、sizeが10を超えると、終了コード3で失敗するコンテナを実行し、continueOnで失敗を吸収し、そのあとのタスクreportが実行されます。チャンクbのノードはFailed、ワークフローとreportはSucceededである必要があります。
retryStrategyは、同じことをもう一度試みることで、continueOnは、失敗を認めたうえで、次に進むことです。ファンアウトの1つの要素が失敗したときに、DAGの依存するタスクがどうなるかを、比べてみてください。
人が承認するまで止まって待つ
/root/capa-data/approve.yamlに、stepsで、split → suspendテンプレートapprove → コンテナテンプレートpublishの順序のWorkflow(generateNameapprove-、ラベルcapa-data: approve)を作成して、--waitなしで提出してください。ワークフローがapproveで止まったことを確認し、10秒以上待ってから、argo resumeで再開して、Succeededで終わらせてください。
suspendテンプレートにdurationを指定しなければ、人が再開するまで待ちます。止まっている間の、ワークフローのphaseとapproveノードのphaseを確認してください。採点は、approveノードが留まった時間を見ます。
2つのワークフローが、同じ倉庫に同時に書き込めないようにする
argoネームスペースに、ConfigMapcapa-data-locks(キーwarehouse、値"1")を作成し、/root/capa-data/locked.yamlに、spec.synchronizationでそのキーをセマフォとして使うWorkflow(generateNamelocked-、ラベルcapa-data: locked、コンテナがsleep 8を実行)を作成してください。同じファイルで、ワークフロー2つを続けて提出して、どちらもSucceededで終わらせます。2つのワークフローの実行区間は、重ならない必要があります。
parallelismは、1つのワークフローの中の上限で、セマフォは、ワークフローの間で共有されるロックです。2つ目のワークフローが待っている間に、status.synchronizationとmessageを確認してください。このバージョン(v4.1.3)は、古い単数形のフィールドsemaphoreを、知らないフィールドとして拒否し、リストのフィールド(semaphores・mutexes)だけを受け付けます(実測)。