すべてのプロダクト
Search
ドキュメントセンター

Container Service for Kubernetes:Python SDK を使用した大規模な Argo Workflows の構築

最終更新日:Jun 19, 2026

Argo Workflows は、スケジュールされたタスク、機械学習、ETL に広く使用されている強力なワークフロー管理ツールです。しかし、YAML でワークフローを定義するには学習コストが高いです。Hera Python SDK は、シンプルな代替手段を提供します。Hera を使用すると、ユーザーは Python でワークフローを構築できます。これにより、複雑なタスクをサポートし、テストを簡素化し、Python エコシステムとシームレスに統合できるため、複雑なワークフローの設計がはるかに容易になります。本トピックでは、Python SDK を使用して大規模な Argo Workflows を構築する方法について説明します。

背景情報

Argo Workflows は、Kubernetes 環境向けに特別に設計されたオープンソースのワークフロー管理ツールです。複雑なワークフローオーケストレーションに重点を置き、ユーザーは一連のタスクを定義し、その実行順序と依存関係を柔軟に設定できます。Argo Workflows は、高度にカスタマイズされた自動化ワークフローを効率的に構築および管理するのに役立ちます。

Argo Workflows には、スケジュールされたタスク、機械学習、シミュレーションコンピューティング、科学計算、ETL、モデルトレーニング、CI/CD など、幅広いユースケースがあります。主に YAML を使用してワークフローを定義しますが、これは明確さとシンプルさを目的とした設計上の選択です。しかし、YAML に慣れていないユーザーにとって、特に複雑なワークフローの場合、その厳密なインデントと階層構造は学習コストが高くなる可能性があります。

image

Hera は、Argo Workflows を構築および送信するために設計された Python SDK フレームワークです。ワークフローの作成と送信を簡素化します。データサイエンティストにとって、Python を使用することは、一般的なプラクティスと一致しており、YAML の課題を克服するのに役立ちます。

作成方法の比較

YAML

Hera

シンプルさ

高い

高い、コード行数が少ない

複雑なワークフローの作成

難しい

簡単、YAML の構文エラーの可能性を効果的に回避できます

Python エコシステムの統合

難しい

簡単、豊富な Python ライブラリにアクセス可能

テスト容易性

難しい、構文エラーが発生しやすい

簡単、テストフレームワークを使用して、コードの品質と保守性を向上させることができます。

Hera は、Python エコシステムと Argo Workflows フレームワークを接続し、ワークフロー設計をより直感的にします。YAML の複雑さなしに大規模なタスクオーケストレーションを可能にし、データサイエンティストやエンジニアが好みの Python 環境で作業できるようにします。これにより、機械学習ワークフローの構築と最適化がシームレスかつ効率的になり、アイデアからデプロイまでのイテレーションサイクルが加速します。以下の例では Hera を使用します。

ステップ1:クラスターの作成とトークンの取得

  1. Argo Workflows クラスターを作成し、次にArgo Server を有効にしてワークフローコンソールにアクセスします。

  2. クラスタートークンを作成します。

    kubectl create token default -n default

ステップ2:Hera を使用したワークフローの送信

  1. Hera をインストールします。

    pip install hera
  2. ワークフローを作成して送信します。

    単純な DAG ダイヤモンド

    Argo Workflows では、有向非巡回グラフ (DAG) が複雑なタスクの依存関係を定義するためによく使用されます。ダイヤモンドは、タスクが分岐してから収束する一般的なパターンです。この構造は、結果が共通の下流タスクに集約される並列処理に効果的です。次の例は、Hera を使用してダイヤモンド構造のワークフローを定義する方法を示しています。タスク A が最初に実行され、次に 2 つの並列タスク B と C が実行されます。最後のタスク D は、B と C の両方が終了した後に実行され、ワークフローが完了します。

    1. simpleDAG.py という名前のファイルを作成し、次の内容を記述します。

      # 必要なパッケージをインポートします。
      from hera.workflows import DAG, Workflow, script
      from hera.shared import global_config
      import urllib3
      urllib3.disable_warnings()
      # ホストアドレスとトークンを設定します。
      global_config.host = "https://{{argo_server_IP}}:2746"
      global_config.token = "abcdefgxxxxxx"  # 取得したトークンに置き換えます。
      global_config.verify_ssl = False
      # @script デコレーターは、Python 関数をほぼネイティブにオーケストレーションできる Hera の主要な機能です。
      # これにより、Workflow や Steps コンテキストなどの Hera コンテキストマネージャー内で装飾された関数を呼び出すことができます。
      # この関数は Hera コンテキスト外でも正常に実行されるため、ユニットテストを作成できます。
      # この例では、入力メッセージを出力します。
      @script(image="mirrors-ssl.aliyuncs.com/python:3.10")
      def echo(message: str):
          print(message)
      # Workflow は Argo の主要なリソースであり、Hera の重要なクラスです。 テンプレートを保存し、エントリポイントを設定して実行します。
      with Workflow(
          generate_name="dag-diamond-",
          entrypoint="diamond",
          namespace="default",
      ) as w:
          with DAG(name="diamond"):
              A = echo(name="A", arguments={"message": "A"})  # テンプレートを構築します。
              B = echo(name="B", arguments={"message": "B"})
              C = echo(name="C", arguments={"message": "C"})
              D = echo(name="D", arguments={"message": "D"})
              A >> [B, C] >> D      # 依存関係を定義します:タスクBとCはAに依存し、タスクDはBとCに依存します。
      # ワークフローを作成します。
      w.create()
    2. ワークフローを送信します。

      python simpleDAG.py
    3. ワークフローの実行後、タスクの DAG とその結果を [ワークフローコンソール] で表示できます。

      ワークフローの例 dag-diamond-g9v45 は、ダイヤモンド型の DAG トポロジーを示しています。最上位ノード A が完了し、次にノード BC が並行して実行され、最後にノード D で収束します。すべてのノードは正常に実行されたとマークされます。

    MapReduce

    Argo Workflows では、DAG テンプレートを使用して MapReduce スタイルのデータ処理を実装し、map フェーズと reduce フェーズをシミュレートできます。次の例は、Hera を使用して単純な MapReduce ワークフローを構築する方法を示しています。このワークフローは、タスクを複数の並列 map タスクに分割し、その結果を最後の reduce タスクで集約します。各ステップは Python 関数であるため、Python エコシステムとの統合が容易になります。

    1. アーティファクトを設定します。

    2. map-reduce.py という名前のファイルを作成し、次の内容を記述します。

      コード

      from hera.workflows import DAG, Artifact, NoneArchiveStrategy, Parameter, OSSArtifact, Workflow, script
      from hera.shared import global_config
      import urllib3
      urllib3.disable_warnings()
      # ホストアドレスを設定します。
      global_config.host = "https://{{argo_server_IP}}:2746"
      global_config.token = "abcdefgxxxxxx"  # 取得したトークンに置き換えます。
      global_config.verify_ssl = False
      # @script デコレーターを使用する場合、image、inputs、outputs、resources などのパラメーターを渡します。
      @script(
          image="mirrors-ssl.aliyuncs.com/python:alpine3.6",
          inputs=Parameter(name="num_parts"),
          outputs=OSSArtifact(name="parts", path="/mnt/out", archive=NoneArchiveStrategy(), key="{{workflow.name}}/parts"),
      )
      def split(num_parts: int) -> None:  # この関数は、num_parts 入力パラメーターに基づいて複数の出力ファイルを作成します。 各ファイルに 'foo' フィールドとパート番号を書き込みます。
          import json
          import os
          import sys
          os.mkdir("/mnt/out")
          part_ids = list(map(lambda x: str(x), range(num_parts)))
          for i, part_id in enumerate(part_ids, start=1):
              with open("/mnt/out/" + part_id + ".json", "w") as f:
                  json.dump({"foo": i}, f)
          json.dump(part_ids, sys.stdout)
      # @script デコレーターで image、inputs、outputs を定義します。
      @script(
          image="mirrors-ssl.aliyuncs.com/python:alpine3.6",
          inputs=[Parameter(name="part_id", value="0"), Artifact(name="part", path="/mnt/in/part.json"),],
          outputs=OSSArtifact(
              name="part",
              path="/mnt/out/part.json",
              archive=NoneArchiveStrategy(),
              key="{{workflow.name}}/results/{{inputs.parameters.part_id}}.json",
          ),
      )
      def map_() -> None:  # この関数は、入力ファイルから 'foo' の値を読み取り、それを2倍にして、結果を新しい出力ファイルの 'bar' フィールドに書き込みます。
          import json
          import os
          os.mkdir("/mnt/out")
          with open("/mnt/in/part.json") as f:
              part = json.load(f)
          with open("/mnt/out/part.json", "w") as f:
              json.dump({"bar": part["foo"] * 2}, f)
      # @script デコレーターで image、inputs、outputs、resources を定義します。
      @script(
          image="mirrors-ssl.aliyuncs.com/python:alpine3.6",
          inputs=OSSArtifact(name="results", path="/mnt/in", key="{{workflow.name}}/results"),
          outputs=OSSArtifact(
              name="total", path="/mnt/out/total.json", archive=NoneArchiveStrategy(), key="{{workflow.name}}/total.json"
          ),
      )
      def reduce() -> None:   # この関数は、すべての map タスクからの 'bar' 値の合計を計算します。
          import json
          import os
          os.mkdir("/mnt/out")
          total = 0
          for f in list(map(lambda x: open("/mnt/in/" + x), os.listdir("/mnt/in"))):
              result = json.load(f)
              total = total + result["bar"]
          with open("/mnt/out/total.json", "w") as f:
              json.dump({"total": total}, f)
      # ワークフローを構築します。 名前、エントリポイント、namespace、およびグローバルパラメーターを定義します。
      with Workflow(generate_name="map-reduce-", entrypoint="main", namespace="default", arguments=Parameter(name="num_parts", value="4")) as w:
          with DAG(name="main"):
              s = split(arguments=Parameter(name="num_parts", value="{{workflow.parameters.num_parts}}")) # テンプレートを構築します。
              m = map_(
                  with_param=s.result,
                  arguments=[Parameter(name="part_id", value="{{item}}"), OSSArtifact(name="part", key="{{workflow.name}}/parts/{{item}}.json"),],
              )   # パラメーターを渡し、テンプレートを構築します。
              s >> m >> reduce()   # タスクの依存関係を定義します。
      # ワークフローを作成します。
      w.create()
      
    3. ワークフローを送信します。

      python map-reduce.py
    4. ワークフローの実行後、タスクの DAG とその結果を [ワークフローコンソール] で表示できます。[WORKFLOW DETAILS] ページでは、DAG ビューには split ノード、4 つの並列 map ノード、および reduce ノードがすべて正常に実行されたことが (緑色のチェックマークで) 示されます。

関連ドキュメント

  • Hera のドキュメント:

  • YAML デプロイメントの例:

    • YAML を使用して単純なダイヤモンドの例をデプロイするには、「dag-diamond.yaml」をご参照ください。

    • YAML を使用して map-reduce の例をデプロイするには、「map-reduce.yaml」をご参照ください。