TT Lab
はじめる
学ぶ 学習パス コース

CAPA — Argoプロジェクト認定アソシエイト

シャードは 3 つなのに Pod は 1 つしか起動しなかった

TT Labで続きを見る

目標

前のステップが作ったデータで、後ろのステップの形が決まるパイプラインを、Argo Workflowsで実際に実行します。スクリプトの出力でファンアウトし、条件でブランチを分け、散らばった結果を集め、失敗を吸収し、承認を待ち、ワークフローの間にロックをかけることを、ノードの記録で確認します。

なぜ重要なのか

データ処理のパイプラインは、実行する前には、いくつのチャンクを処理するのかがわかりません。入力を分割するステップが動いてはじめて、チャンクの数が決まり、その数だけ並列の作業ができる必要があります。YAMLにリストをあらかじめ書く方式では、これを表現できないため、Argo Workflowsは、あるノードの出力を、次のノードの繰り返しリストとして使うwithParam、要素ごとに評価されるwhen、ファンアウトの出力の自動集計を用意しています。その代わり、並列性が基本なので、共有リソースを使うステップには、上限とロックを自分でかける必要があり、チャンク1つの失敗を、全体の失敗と見なすかどうかも、設計で決める必要があります。試験は、このテンプレートの種類とspecのフィールドが、実行で何を変えるかを問います。

ステップ

  1. /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で出力します。
  2. /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である必要があります。
  3. /root/capa-data/serial.yamlに、ステップ2と同じファンアウトのWorkflowを、generateNameserial-、ラベルcapa-data: serialで作成し、processコンテナがsleep 3を実行し、DAGテンプレートmainにparallelism: 1を置いて、提出してください。3つのprocessのPodの実行区間が、互いに重ならない必要があります。
  4. /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として残る必要があります。
  5. /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の合計である必要があります。
  6. /root/capa-data/tolerate.yamlに、Workflow(generateNametolerate-、ラベルcapa-data: tolerate)を作成して、提出してください。ファンアウトのタスクcheckは、sizeが10を超えると、終了コード3で失敗するコンテナを実行し、continueOnで失敗を吸収し、そのあとのタスクreportが実行されます。チャンクbのノードはFailed、ワークフローとreportはSucceededである必要があります。
  7. /root/capa-data/approve.yamlに、stepsで、split → suspendテンプレートapprove → コンテナテンプレートpublishの順序のWorkflow(generateNameapprove-、ラベルcapa-data: approve)を作成して、--waitなしで提出してください。ワークフローがapproveで止まったことを確認し、10秒以上待ってから、argo resumeで再開して、Succeededで終わらせてください。
  8. argoネームスペースに、ConfigMapcapa-data-locks(キーwarehouse、値"1")を作成し、/root/capa-data/locked.yamlに、spec.synchronizationでそのキーをセマフォとして使うWorkflow(generateNamelocked-、ラベルcapa-data: locked、コンテナがsleep 8を実行)を作成してください。同じファイルで、ワークフロー2つを続けて提出して、どちらもSucceededで終わらせます。2つのワークフローの実行区間は、重ならない必要があります。

参考

スクリプトの標準出力がデータになる

/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)だけを受け付けます(実測)。