概要
dagは「本当に必要な依存関係だけ」を宣言する
steps(配列の配列)は暗黙的に直列性を強制する。依存関係のない処理までstepsで並べると意図せず直列化し所要時間が悪化する。dag+dependenciesで必要な依存だけを明示すれば、Argoのスケジューラが安全な範囲で自動並列化してくれる。
再試行すべき失敗と、しても無駄な失敗を区別する
BigQueryのレート制限やネットワークの瞬断は再試行で解決するが、SQL構文エラーや権限不足は何度再試行しても失敗する。retryPolicy: OnTransientErrorで両者を区別しないと、恒久的エラーの検知がlimit回数分だけ遅延する。
Artifactsはデータ受け渡しを「宣言」に変える
文字列で組み立てたGCSパスはArgo自身にとってはただの文字列で、リネージ追跡もタイポ検知もできない。outputs.artifacts/inputs.artifactsで受け渡しを宣言すると、Argo UIで生成・消費関係を可視化できる。
onExitは「失敗を握り潰さない」ための最終防衛線
通常のstepとして通知を実装すると、そのstepより前で失敗した場合に通知自体が実行されない構造的な穴が生まれる。spec.onExitは成功・失敗を問わず必ず実行されるため、「気づけない失敗」を構造的に排除できる。
問題
MOps 販促システムチームでは、日次で「キャンペーン成果集計バッチ」(campaign-performance-daily)を GKE Autopilot 上の Argo Workflows で運用している。処理内容はextract-orders(注文イベント抽出)→extract-coupons(クーポン使用ログ抽出、ordersとは独立で依存関係なし)→transform(成果Mart計算)→load(BigQueryへロード+Slack通知)の4段階。
制約・前提条件
- Argo Workflows 3.5系、GKE Autopilot(Kubernetes 1.30+)上で動作
extract-ordersとextract-couponsは異なるBigQueryテーブルを対象とし、互いに依存しない(並列実行可能)bq-sa-key-plaintextの鍵ファイルは廃止し、Workload Identity Federationで認証したい(GCP側のIAMバインディングは設定済み)- Slack通知は成功時だけでなく失敗時にも必ず飛ばしたい(現状は成功パスのstepとしてしか実装されておらず、失敗時は通知自体が実行されない)
- BigQuery APIの一時的なレート制限(429)やGCSの一時的なネットワークエラーは自動リトライで吸収したいが、SQL構文エラーや権限エラーなど恒久的なエラーは即座に失敗させたい
target-dateパラメータはYYYY-MM-DD形式であることを起動時点で検証したい
悪いコード (Before)
apiVersion: argoproj.io/v1alpha1
kind: WorkflowTemplate
metadata:
name: campaign-performance-daily
spec:
entrypoint: main
arguments:
parameters:
- name: target-date
templates:
- name: main
steps:
# 問題①: 依存関係のないextract-orders/extract-couponsまで直列化
- - name: extract-orders
template: extract-orders
- - name: extract-coupons
template: extract-coupons
- - name: transform
template: transform
- - name: load
template: load
- - name: notify
template: notify
- name: extract-orders
container:
image: gcr.io/mops-prod/bq-extractor:latest
command: [python, extract.py]
args:
- "--table=orders"
- "--date={{workflow.parameters.target-date}}"
# 問題③: GCSパスを文字列連結でハードコード
- "--out=gs://mops-tmp/orders/{{workflow.parameters.target-date}}.parquet"
env:
# 問題④: 長期鍵ファイルをSecretからマウントして参照
- name: GOOGLE_APPLICATION_CREDENTIALS
value: /secrets/sa-key.json
volumeMounts:
- name: sa-key
mountPath: /secrets
volumes:
- name: sa-key
secret:
secretName: bq-sa-key-plaintext
# 問題②: retryStrategy未設定
# 問題⑤: activeDeadlineSeconds未設定
- name: extract-coupons
container:
image: gcr.io/mops-prod/bq-extractor:latest
command: [python, extract.py]
args:
- "--table=coupon_usage"
- "--date={{workflow.parameters.target-date}}"
- "--out=gs://mops-tmp/coupons/{{workflow.parameters.target-date}}.parquet"
- name: transform
container:
image: gcr.io/mops-prod/campaign-transform:latest
command: [python, transform.py]
args:
- "--orders=gs://mops-tmp/orders/{{workflow.parameters.target-date}}.parquet"
- "--coupons=gs://mops-tmp/coupons/{{workflow.parameters.target-date}}.parquet"
- "--out=gs://mops-tmp/mart/{{workflow.parameters.target-date}}.parquet"
- name: load
container:
image: gcr.io/mops-prod/bq-loader:latest
command: [python, load.py]
args:
- "--in=gs://mops-tmp/mart/{{workflow.parameters.target-date}}.parquet"
- "--table=campaign_performance_daily"
- name: notify
# 問題⑥: 成功パスの最後のstepとしてしか通知しない
# → これより前段で失敗すると通知自体が実行されない(今回の事故の直接原因)
# 問題⑦: target-dateの形式検証が一切なく不正な値のまま課金対象のBQジョブが走る
container:
image: curlimages/curl
command: [sh, -c]
args:
- "curl -X POST -d '{\"text\":\"campaign-performance-daily completed\"}' https://hooks.slack.com/services/T000/B000/XXXXXXXXXXXX"
ヒント(段階的開示)
ヒント1 — 方向性
steps(配列の配列=直列のブロック、内側の配列は並列)を使っているが、実質すべてのステップが1つずつのブロックになっており、依存関係のないextract-ordersとextract-couponsまで直列化されている。dagテンプレートでdependenciesを宣言すればスケジューラが自動的に並列実行してくれる。(2) 失敗の扱いが「全か無か」になっている — リトライが一切ないため一時的な障害でも即失敗し、かつ失敗時にexitHandlerが無いため誰にも気づかれない。(3) データ受け渡しが「命名規約という名の口約束」になっている — GCSパスを文字列で毎回組み立てており、Argo自身がこのファイルの生成・消費関係を認識できていない。(4) 認証情報の運用がCLAUDE.mdのセキュリティ原則に反する — 長期鍵をSecretとしてマウントしている。
「このワークフローが3ヶ月後、深夜2時にBigQueryの一時障害で失敗したとき、誰が・どうやって気づくか?」を自問すると、指摘の半分(リトライ・exitHandler)が自然に出てくる。
ヒント2 — アプローチ
- 問題①:
stepsで直列実行 →dagテンプレートに変更し、extract-orders/extract-couponsは共通の依存元(バリデーションstep)にのみ依存させ、両者の間には依存関係を張らない - 問題②: リトライなし → 各templateに
retryStrategyを設定し、retryPolicy: OnTransientError+backoff(指数バックオフ)+limit(上限回数)を組み合わせる - 問題③: GCSパスのハードコード →
outputs.artifacts/inputs.artifactsを使い、各templateはローカルパスだけを意識すればArgoが裏側でアップロード/ダウンロードを自動化してくれる形にする - 問題④: SAキーファイルのマウント →
volumes/GOOGLE_APPLICATION_CREDENTIALSを削除し、PodにserviceAccountNameでWorkload Identityバインド済みKSAを指定する - 問題⑤: タイムアウトなし → 各templateに
activeDeadlineSecondsを設定する - 問題⑥: 失敗時に気づけない →
spec.onExitにnotify-exitテンプレートを指定する。{{workflow.status}}を参照でき、成功・失敗どちらでも必ず実行される - 問題⑦:
target-dateの未検証 → DAGの先頭ノードとしてvalidate-dateテンプレートを置き、他の全ステップがこれに依存する形にする
ヒント3 — コードの骨格
apiVersion: argoproj.io/v1alpha1
kind: WorkflowTemplate
metadata:
name: campaign-performance-daily
spec:
entrypoint: main
onExit: notify-exit # 成功/失敗いずれでも必ず呼ばれる
serviceAccountName: campaign-pipeline-ksa # Workload Identityでバインド済みKSA、鍵ファイル不要
arguments:
parameters:
- name: target-date
templates:
- name: main
dag:
tasks:
- name: validate-date
template: validate-date
- name: extract-orders
template: extract-orders
dependencies: [validate-date]
- name: extract-coupons
template: extract-coupons
dependencies: [validate-date] # extract-ordersとは依存しない→並列実行される
- name: transform
template: transform
dependencies: [extract-orders, extract-coupons]
arguments:
artifacts:
- name: orders
from: "{{tasks.extract-orders.outputs.artifacts.orders-parquet}}"
- name: coupons
from: "{{tasks.extract-coupons.outputs.artifacts.coupons-parquet}}"
- name: load
template: load
dependencies: [transform]
arguments:
artifacts:
- name: mart
from: "{{tasks.transform.outputs.artifacts.mart-parquet}}"
- name: extract-orders
outputs:
artifacts:
- name: orders-parquet
path: /work/orders.parquet # Argoが自動でArtifact Repositoryへアップロード
retryStrategy:
limit: 3
retryPolicy: OnTransientError
backoff: { duration: 30s, factor: 2, maxDuration: 5m }
activeDeadlineSeconds: 900
# container定義(image, command, args, resources)は省略
- name: notify-exit
container:
image: curlimages/curl:8.7.1
command: [sh, -c]
args:
- >
STATUS="{{workflow.status}}";
curl -sf -X POST -H 'Content-Type: application/json'
-d "{\"text\":\"campaign-performance-daily [$STATUS]\"}"
"$SLACK_WEBHOOK_URL"
envFrom:
- secretRef:
name: slack-webhook # Webhook URLはコード直書きせずSecretから注入
問題点分析(7点)
| # | 問題点 | 分類 | 改善方法 |
|---|---|---|---|
| 1 | 依存関係のないextract-orders/extract-couponsをstepsで直列実行 | 実行グラフ | dagテンプレート+dependenciesで必要な依存だけ宣言、自動並列化 |
| 2 | retryStrategy皆無で一時的な429エラーも即失敗 | 耐障害性 | retryPolicy=OnTransientError+指数バックオフ+limit |
| 3 | GCSパスの文字列連結でArgo Artifacts未使用 | データ受け渡し | outputs.artifacts/inputs.artifactsで宣言的に受け渡し |
| 4 | SAキーファイルをSecretマウントして参照 | 認証情報運用 | Workload Identity Federationで鍵レス化 |
| 5 | activeDeadlineSeconds未設定でハングが無期限ブロック | タイムアウト設計 | 各templateにactiveDeadlineSecondsを設定 |
| 6 | onExit未設定で失敗時に通知されず1日気づかれない | 可観測性 | spec.onExitで成功/失敗いずれも必達通知 |
| 7 | target-date未検証のまま課金対象のBQジョブ実行 | 入力検証 | validate-dateをDAG先頭ノード化、全stepが依存 |
構成図 — Bad vs Good(SVG)
模範解答
templates:
- name: main
steps:
- - name: extract-orders
template: extract-orders
- - name: extract-coupons
template: extract-coupons
- - name: transform
template: transform
- - name: load
template: load
- - name: notify
template: notify
apiVersion: argoproj.io/v1alpha1
kind: WorkflowTemplate
metadata:
name: campaign-performance-daily
spec:
entrypoint: main
onExit: notify-exit # 修正⑥: 成功/失敗どちらでも必ず呼ばれる
serviceAccountName: campaign-pipeline-ksa # 修正④: Workload IdentityでバインドされたKSA
arguments:
parameters:
- name: target-date
templates:
- name: main
dag:
tasks:
- name: validate-date
template: validate-date
arguments:
parameters:
- name: target-date
value: "{{workflow.parameters.target-date}}"
- name: extract-orders
template: extract-orders
dependencies: [validate-date] # 修正⑦: バリデーション後にのみ実行
arguments:
parameters:
- name: target-date
value: "{{workflow.parameters.target-date}}"
- name: extract-coupons
template: extract-coupons
dependencies: [validate-date] # 修正①: ordersと依存なし→両者は並列実行される
arguments:
parameters:
- name: target-date
value: "{{workflow.parameters.target-date}}"
- name: transform
template: transform
dependencies: [extract-orders, extract-coupons] # 両方の完了を待つ
arguments:
artifacts:
- name: orders
from: "{{tasks.extract-orders.outputs.artifacts.orders-parquet}}"
- name: coupons
from: "{{tasks.extract-coupons.outputs.artifacts.coupons-parquet}}"
- name: load
template: load
dependencies: [transform]
arguments:
artifacts:
- name: mart
from: "{{tasks.transform.outputs.artifacts.mart-parquet}}"
- name: validate-date
inputs:
parameters:
- name: target-date
container:
image: gcr.io/mops-prod/date-validator:1.4
command: [python, validate.py]
args: ["--date", "{{inputs.parameters.target-date}}"]
activeDeadlineSeconds: 60
container:
args:
- "--out=gs://mops-tmp/orders/{{workflow.parameters.target-date}}.parquet"
env:
- name: GOOGLE_APPLICATION_CREDENTIALS
value: /secrets/sa-key.json
volumeMounts:
- name: sa-key
mountPath: /secrets
volumes:
- name: sa-key
secret:
secretName: bq-sa-key-plaintext
# retryStrategy / activeDeadlineSeconds も無し
- name: extract-orders
inputs:
parameters:
- name: target-date
outputs:
artifacts:
- name: orders-parquet
path: /work/orders.parquet # 修正③: Argoが自動でGCSへアップロード
container:
image: gcr.io/mops-prod/bq-extractor:1.4
command: [python, extract.py]
args:
- "--table=orders"
- "--date={{inputs.parameters.target-date}}"
- "--out=/work/orders.parquet" # ローカルパスのみ意識すればよい
# GOOGLE_APPLICATION_CREDENTIALS明示指定は不要
# → serviceAccountName(Workload Identity)が自動適用される
resources:
requests: { cpu: "500m", memory: "1Gi" }
limits: { cpu: "1", memory: "2Gi" }
activeDeadlineSeconds: 900 # 修正⑤: 15分でタイムアウト
retryStrategy: # 修正②: 一時的な障害のみ再試行
limit: 3
retryPolicy: OnTransientError # 429やネットワーク一時障害のみ再試行
backoff:
duration: 30s
factor: 2
maxDuration: 5m
# --- 成果Mart計算: Artifactsで入出力を宣言的に受け渡し ---
- name: transform
inputs:
artifacts:
- name: orders
path: /work/orders.parquet # Argoが対応するArtifactを自動でダウンロードして配置
- name: coupons
path: /work/coupons.parquet
outputs:
artifacts:
- name: mart-parquet
path: /work/mart.parquet
container:
image: gcr.io/mops-prod/campaign-transform:1.4
command: [python, transform.py]
args:
- "--orders=/work/orders.parquet"
- "--coupons=/work/coupons.parquet"
- "--out=/work/mart.parquet"
resources:
requests: { cpu: "1", memory: "2Gi" }
limits: { cpu: "2", memory: "4Gi" }
activeDeadlineSeconds: 1200
retryStrategy:
limit: 2
retryPolicy: OnTransientError
# --- BigQueryへロード ---
- name: load
inputs:
artifacts:
- name: mart
path: /work/mart.parquet
container:
image: gcr.io/mops-prod/bq-loader:1.4
command: [python, load.py]
args:
- "--in=/work/mart.parquet"
- "--table=campaign_performance_daily"
resources:
requests: { cpu: "500m", memory: "1Gi" }
limits: { cpu: "1", memory: "2Gi" }
activeDeadlineSeconds: 600
retryStrategy:
limit: 3
retryPolicy: OnTransientError
backoff: { duration: 30s, factor: 2, maxDuration: 5m }
- name: notify
container:
image: curlimages/curl
command: [sh, -c]
args:
- "curl -X POST -d '{\"text\":\"completed\"}' https://hooks.slack.com/..."
# extract-orders等が先に失敗すると、このstep自体が実行されない
spec:
onExit: notify-exit # ワークフロー全体の外側で必ず実行される
templates:
- name: notify-exit
container:
image: curlimages/curl:8.7.1
command: [sh, -c]
args:
- >
STATUS="{{workflow.status}}";
MSG="campaign-performance-daily [$STATUS] date={{workflow.parameters.target-date}}";
curl -sf -X POST -H 'Content-Type: application/json'
-d "{\"text\":\"$MSG\"}"
"$SLACK_WEBHOOK_URL"
envFrom:
- secretRef:
name: slack-webhook # Webhook URLをYAMLに直書きせずSecretから注入
activeDeadlineSeconds: 30
起動: argo submit --from workflowtemplate/campaign-performance-daily -p target-date=2026-08-16
Before: extract-orders → extract-coupons → transform → load → notify を完全直列実行(約12分)
extract-orders でBigQuery 429エラー発生 → リトライなしで即Failed → ワークフロー全体が停止
notify stepは成功パスの最後にしかないため実行されず、誰も失敗に気づかない
After : validate-date(date形式チェック)→ extract-orders / extract-coupons が並列実行(約7分に短縮)
→ transform → load
extract-orders でBigQuery 429エラー発生
→ retryStrategy(limit=3, backoff 30s→1m→2m)で自動リトライ → 2回目で成功、パイプライン継続
仮に3回とも失敗した場合:
→ タスクFailed → ワークフロー全体がFailedへ遷移
→ onExit: notify-exit が必ず実行され、Slackに "[Failed] date=2026-08-16" が通知される
→ SREが即座に気づいて手動再実行できる
| 問題 | 修正内容 | 効果 |
|---|---|---|
| ① 依存なしを直列化 | dag+dependenciesで必要な依存のみ宣言 | 並列実行され所要時間が短縮 |
| ② retryStrategy皆無 | OnTransientError+指数バックオフ+limit | 一時障害を自動吸収、恒久障害は即失敗 |
| ③ GCSパス文字列連結 | outputs.artifacts/inputs.artifactsで受け渡し | リネージ追跡可能、タイポ耐性が上がる |
| ④ SAキーファイルマウント | Workload Identity Federationで鍵レス化 | 長期鍵の漏洩リスク・運用コストがゼロに |
| ⑤ activeDeadlineSeconds無し | 各templateにタイムアウトを設定 | ハングによる無限コスト発生を防止 |
| ⑥ onExit未設定 | spec.onExitで成功/失敗いずれも必達通知 | 失敗を握り潰さず必ず気づける |
| ⑦ target-date未検証 | validate-dateをDAG先頭ノード化 | 不正な値で高コストなBQジョブが走らない |
ポイント解説
stepsは「このブロックは前のブロックの後」という直列性を暗黙的に強制する。依存関係のない処理までstepsで並べると意図せず直列化してしまい所要時間もリソース効率も悪化する。dag+dependenciesは「本当に必要な依存関係だけ」を明示するため、スケジューラが安全な範囲で最大限並列化してくれる。limit回数分だけ遅延し、無駄なコストと待ち時間を生む。outputs.artifacts/inputs.artifactsで受け渡しを宣言すると、Argo UI上で「どのタスクの出力が、どのタスクの入力になったか」を可視化でき、命名規約に依存しない安全な受け渡しになる。onExitは成功・失敗を問わず必ず実行されるため、「気づけない失敗」を構造的に排除できる。実務への応用
既存のcampaign-performance-dailyワークフローを今回の設計に置き換える際、dag化による並列実行短縮効果は「セール当日の集計遅延」を直接改善できる(08-12で扱ったBI Engine/Materialized View施策と組み合わせれば、抽出〜可視化までのEnd-to-Endレイテンシをさらに縮められる)。
retryPolicy: OnTransientErrorの考え方は、Cloud Run上のAPI(08-15で扱ったcoupon-service等)のリトライ設計にも転用できる。「一時的なエラーだけ自動リトライし、恒久的なエラーは即座にアラートへ回す」という区別は、Argo Workflowsに限らずAPIクライアントのRetry-Afterヘッダー処理やCircuit Breaker設計とも共通する原則。
onExitによるSlack通知は、他のバッチ(クーポン発行バッチ、レコメンド特徴量生成バッチ)にもテンプレート化して展開できる。WorkflowTemplateを共通のonExit定義付きベーステンプレートとして整備すれば、チーム全体のバッチで「失敗に気づけない」事故を構造的に防止できる。
今日のまとめ
dagによる正確な依存関係表現、retryStrategyによる一時障害と恒久障害の区別、onExitによる失敗の必達通知、Artifactsによる宣言的なデータ受け渡し、Workload Identityによる鍵レス認証——これら5つはいずれも「パイプラインが3ヶ月後に深夜落ちても、翌朝までに人間が気づいて対処できるか」という1つの問いから導かれる設計原則である。
次のステップ
- 発展問題:
extract-orders/extract-couponsのBigQueryジョブが恒久的エラー(権限不足)で失敗した場合に、retryStrategyのlimitを無駄に消費せず即座に失敗と判定する仕組み(exit codeによるOnTransientErrorの判定ロジック)をbq-extractorイメージ側にどう実装するか設計する。 - 発展問題:
WorkflowTemplateをCronWorkflow化し、日次実行のスケジューリング・タイムゾーン(JST)・実行済みチェック(重複実行防止)を追加する。 - 参考: Argo Workflows公式ドキュメント(DAG Templates / Retry Strategy / Artifacts / Exit Handlers)、GCP Workload Identity Federation for GKEドキュメント