KFP(Kubeflow Pipelines)は、機械学習の前処理・学習・評価・デプロイをコンテナの連なり(DAG)として定義し、Kubernetes上で実行・記録する仕組みです。v2ではPythonのデコレーターでパイプラインを書き、Argo Workflowsに依存しない中間表現(IR YAML)へコンパイルする形に変わりました。この記事では、2026年10月時点の最新SDK(kfp 2.17.0)で実際に動かしたコードをもとに、書き方・v1からの移行・実行先の選び方をまとめます。Kubeflow全体の構成や導入経路はKubeflowの全体像とPipelinesの使い方で扱っています。
まとめ:KFP v2を使い始める前に押さえる要点
- 版:SDKはkfp 2.17.0(2026年7月9日公開・Python 3.9以上)、バックエンドは2.17.2(2026年9月4日)が最新です。
- 書き方:
@dsl.componentで部品を作り、@dsl.pipelineでつなぎ、IR YAMLへコンパイルします。v1のContainerOpやfunc_to_container_opはv2のSDKにありません。 - 受け渡し:数値や文字列はParameter、ファイルやモデルはArtifactとして型注釈で区別します。
- キャッシュは既定で有効です。外部データを読むタスクは明示的に切らないと、前回の結果が再利用されます。
- 手元実行:
kfp localはSDK 2.16.0以降で条件分岐とループも動きます(公式ページの「未対応」記述は古いままです)。 - v1の扱い:v1はバックエンド2.5から非推奨です。削除時期は未確定ですが、新規コードはv2で書きます。
- Google Cloudで動かすなら、Agent Platform Pipelines(旧Vertex AI Pipelines)が同じIR YAMLをそのまま受け付けます。
KFP v2の仕組み:IR YAMLとバックエンドの役割分担
v1のSDKはパイプラインをArgo WorkflowsのYAMLへ直接変換していました。v2ではSDKが出力するのは実行基盤に依存しないIR YAML(PipelineSpec)で、それをどこで動かすかは実行側が決めます。同じIR YAMLを、セルフホストのKFPバックエンドにもAgent Platform Pipelinesにも投入できるのはこのためです。
ただし、OSSのKFPバックエンドは今も内部の標準オーケストレーターとしてArgo Workflowsを使っています。v2で切り離されたのは「ユーザーが書く定義とArgoの関係」であり、クラスタからArgoが消えたわけではありません。各タスクの実行時には、入力を解決してPodの設定を組み立てるdriverと、アーティファクトを転送してPython関数を呼び出すlauncherが動きます。
| パッケージ/リポジトリ | 最新版 | 公開日 |
|---|---|---|
| kfp(SDK) | 2.17.0 | 2026-07-09 |
| kfp-kubernetes(拡張) | 2.17.0 | 2026-07-09 |
| KFPバックエンド | 2.17.2 | 2026-09-04 |
| kubeflow/manifests | 26.03.1 | 2026-06-15 |
| google-cloud-aiplatform | 2.3.0 | 2026-09-30 |
直近のバックエンドでは、運用に響く変更が続いています。2.15.0では標準のオブジェクトストアがMinIOからSeaweedFSに替わりました(MinIOを含むS3互換ストアは引き続き使えます)。2.17.0では、2.16.1以前のAPIサーバーにあった認証済み利用者によるSSRFの脆弱性(GHSA-f7vj-6669-qvv4)が修正されています。2.16系以前を動かしているなら、2.17系へ上げる理由になります。
コンポーネント3種の書き分け:Lightweight・Containerized Python・Container
v2のコンポーネントは書き方で3種類に分かれます。迷ったら、依存の少ないPython処理はLightweight、社内ライブラリを抱えるならContainerized Python、Python以外や既存イメージはContainerと決めると判断が速くなります。
| 種類 | 書き方 | イメージ | 向く処理 |
|---|---|---|---|
| Lightweight Python | @dsl.component | base_imageで実行時にコードを注入 | 依存が少ない前処理・集計 |
| Containerized Python | @dsl.component+target_image | kfp component buildで作成 | 複数ファイル・社内ライブラリ |
| Container | @dsl.container_component | 既存イメージをそのまま指定 | シェル・Go・バイナリ |
Lightweightは関数本体を実行時にコンテナへ渡すため、関数の外にあるimportやヘルパー関数を参照できません。SDK 2.15.0からはadditional_funcs引数で共通関数を埋め込めるようになりましたが、関数が増えてきたらContainerized Pythonへ移るのが素直です。Containerized Pythonはkfp component buildでソースと依存を新しいイメージに固め、--push-imageを付けるとレジストリへのpushまで行います(--no-push-imageならローカルビルドだけです)。Containerはdsl.ContainerSpecでimage・command・argsを直接書く方式で、v1のコンポーネントYAMLやContainerOpに相当します。
ParameterとArtifact:型注釈で決まる入出力
v2では、関数の引数と戻り値の型注釈がそのまま入出力の定義になります。int・float・str・list・dict・boolはParameterとして値が直接渡され、Input[Dataset]やOutput[Model]のように書いたものはArtifactとしてオブジェクトストレージ上のファイルで渡されます。Artifactは.pathで読み書きし、.metadataに学習率などの付帯情報を残せます。
次のコードは、データ作成→3通りの係数によるスコア計算→最大スコアの選択→条件分岐を示す簡略例です。実際のモデル学習・精度評価・デプロイは行わず、ファイルの受け渡しと制御フローを確認します。kfp 2.17.0で実行し、3回の計算のスコア(0.5・2.5・5.0)が集約されて最大値5.0が選ばれ、「deploy」側の分岐が動くことを確認しています。
from kfp import dsl, compiler
from kfp.dsl import Input, Output, Dataset, Model, Metrics
@dsl.component(base_image="python:3.11")
def make_data(n: int, data: Output[Dataset]) -> int:
with open(data.path, "w") as f:
for i in range(n):
f.write(f"{i},{i*2}\n")
return n
@dsl.component(base_image="python:3.11")
def train(data: Input[Dataset], lr: float,
model: Output[Model], metrics: Output[Metrics]) -> float:
rows = open(data.path).read().splitlines()
score = round(len(rows) * lr, 3)
with open(model.path, "w") as f:
f.write(f"lr={lr}")
model.metadata["lr"] = lr
metrics.log_metric("score", score)
return score
@dsl.component(base_image="python:3.11")
def pick_best(scores: list) -> float:
return max(scores)
@dsl.component(base_image="python:3.11")
def notify(msg: str):
print(msg)
@dsl.pipeline(name="kfp-v2-demo")
def demo(n: int = 5):
d = make_data(n=n)
with dsl.ParallelFor(items=[0.1, 0.5, 1.0]) as lr:
t = train(data=d.outputs["data"], lr=lr)
best = pick_best(scores=dsl.Collected(t.outputs["Output"]))
with dsl.If(best.output > 3.0):
notify(msg="deploy")
with dsl.Else():
notify(msg="skip")
compiler.Compiler().compile(demo, "demo.yaml")
関数の戻り値はOutputという名前の出力になり、t.outputs["Output"](単一出力なら.output)で参照します。v2は型の検査がv1より厳格で、floatの引数に文字列"0.1"を渡すとコンパイルで止まります。また、コンポーネント呼び出しはキーワード引数が必須です。
制御フローの書き方:ParallelFor・Collected・If・ExitHandler
v2の制御フローはwith文のコンテキストマネージャーで書きます。ParallelForとCollected、ExitHandlerはSDK 2.0.0から、If・Elif・ElseはSDK 2.2.0から使えます。
- dsl.ParallelFor:itemsの要素ごとにタスクを並列実行します。
parallelismで同時実行数の上限を決められ、0は無制限です。 - dsl.Collected:ループ内の出力をリストにまとめ、ループ外のタスクへ渡します(fan-in)。v1でよく使われた「ループ内のタスクに
.after()でつなぐ」書き方は、移行ガイドでは「v2ではコンパイルできない」とされていますが、SDK 2.5.0で対応済みで、2.17.0でもコンパイルが通ることを確認しました。実行順だけを指定する場合は.after()を使い、ループ内の出力をリストとして渡す場合はCollectedを使います。 - dsl.If/Elif/Else:パイプライン入力や上流タスクの出力との比較式で分岐します。旧来の
dsl.Conditionは残っていますが、2.17.0で使うと「dsl.Ifを使え」というDeprecationWarningが出ます。 - dsl.ExitHandler:囲んだタスクが失敗しても、最後に
exit_taskを必ず実行します。後片付けや通知に使います。
ExitHandlerには詰まりやすい制約があります。ExitHandlerの中のタスクに、外側のタスクから依存するとコンパイルがInvalidTopologyException(Illegal task dependency across DSL context managers)で止まります。分岐を後ろにつなぎたい場合は、分岐ごとExitHandlerの内側へ入れてください。なお、SDK 2.17.0からはExitHandlerのグループ全体に.after()で依存させる書き方が通るようになりました。
コンパイルと実行:kfp localからクラスタ投入まで
先のPythonコードをpipe.pyとして保存します。検証例の環境はPython 3.11、kfp 2.17.0です。Kubernetes拡張の例にはkfp-kubernetes 2.17.0も必要です。パイプラインはcompiler.Compiler().compile()か、CLIのkfp dsl compileでIR YAMLにします。ここから先は、どこで動かすかで手順が分かれます。
kfp localで手元実行する手順と効かない設定
kfp localを使うと、クラスタを用意せずにパイプラインのロジックを確かめられます。local.init()のあとにパイプライン関数をそのまま呼ぶだけです。
from kfp import local
local.init(runner=local.SubprocessRunner(use_venv=False))
demo(n=5) # パイプライン関数をそのまま呼ぶと手元で実行される
ランナーは2種類あります。SubprocessRunnerはLightweight Pythonコンポーネント専用で、独自イメージは扱えません。use_venv=Falseなら今のPython環境で動き、既定のTrueではタスクごとに仮想環境を作ってkfpを入れ直すため遅くなります。DockerRunnerは3種のコンポーネントすべてを実行できますが、Dockerが必要です。
公式の手元実行ページには、2026年10月3日時点でも条件分岐とループは手元では動かないと書かれています。しかし、2.16.0のリリースノートでdsl.ConditionとParallelForの手元実行対応が追加されており、上の例も2.17.0のSubprocessRunnerでParallelFor・Collected・If/Elseまで完走しました。ExitHandlerも、内側のタスクを失敗させると後片付けのタスクが実行されたうえでパイプラインが失敗として終わることを確認しています。一方、キャッシュ・リトライ・リソース指定・アフィニティは手元では効きません。リトライの挙動などはクラスタで確かめる必要があります。
キャッシュの既定と切り方
KFPは全タスクの実行キャッシュが既定で有効です。入力とコンポーネント定義が同じなら前回の出力が再利用されるため、「毎回外部APIやDBから最新データを取る」タスクは古い結果のまま進みます。切り方は3段階あります。
| 範囲 | 指定方法 |
|---|---|
| タスク単位 | task.set_caching_options(enable_caching=False) |
| コンパイル時の既定 | KFP_DISABLE_EXECUTION_CACHING_BY_DEFAULT=true(SDK 2.10.0以降) |
| 実行投入時(全タスク) | enable_caching=False(タスク側の設定を上書き) |
# 全タスクのキャッシュを既定で切ってコンパイルする(SDK 2.10.0以降)
KFP_DISABLE_EXECUTION_CACHING_BY_DEFAULT=true kfp dsl compile --py pipe.py --output demo.yaml
# 同じ指定をフラグで書く場合
kfp dsl compile --py pipe.py --output demo.yaml --disable-execution-caching-by-default
環境変数を付けてコンパイルすると、IR YAMLからenableCache: trueが消えることを確認しています。タスク単位で無効にした場合、そのタスクのcachingOptionsは空({})になります。
kfp-kubernetesによるSecret・PVCの設定とタスクのリソース指定
v1のVolumeOpやResourceOpはv2のSDKから外れ、Kubernetes固有の設定は拡張ライブラリのkfp-kubernetesへ移りました。Secret・PVC・imagePullPolicy・一時ボリューム・ノードセレクター・トレラレーション・ラベルとアノテーションをタスクに付けられます。PVCはCreatePVCで作り、mount_pvcでマウントし、DeletePVCで後片付けします。
from kfp import dsl, compiler, kubernetes
@dsl.component(base_image="python:3.11")
def step(x: int) -> int:
return x + 1
@dsl.pipeline(name="k8s-demo")
def p():
s = step(x=1)
s.set_caching_options(enable_caching=False) # このタスクだけキャッシュしない
s.set_cpu_limit("1").set_memory_limit("1G").set_retry(num_retries=2)
kubernetes.use_secret_as_env(
s, secret_name="db-cred", secret_key_to_env={"password": "DB_PASS"})
compiler.Compiler().compile(p, "k8s-demo.yaml")
このコードをコンパイルすると、SecretはIR YAMLの本体ではなくplatforms.kubernetes側にsecretAsEnvとして書き出されます。CPU・メモリ上限とリトライ回数は本体側に入ります。プラットフォーム固有の設定は本体と分けて保持されます。公式ドキュメントがkfp-kubernetesの対応先として挙げているのはOSSのKFPバックエンドだけなので、Agent Platform Pipelinesでも同じ設定が効く前提は置かないでください。
セルフホストKFPへのパイプライン投入
KFPバックエンドへは、SDKのkfp.ClientからIR YAMLを投入します。実験名や投入時のキャッシュ指定もここで渡します。
import kfp
client = kfp.Client(host="http://localhost:8080") # ml-pipeline-ui をポートフォワードした例
client.create_run_from_pipeline_package(
"demo.yaml",
arguments={"n": 100},
experiment_name="kfp-v2-demo",
enable_caching=False, # 指定するとタスク側の設定をすべて上書きする
)
SDKとバックエンドは、使用する機能に対応した組み合わせを選び、両方のバージョンを固定してください。SDK 2.15でコンパイルしたパイプラインが古いバックエンドでunknown field "custom_path"として失敗する不具合があり、SDK 2.15.2で修正されました。コンパイルに使うSDKをバックエンドより新しくしないのが安全です。
v1からv2への移行:書き換え対応表と止まりやすい箇所
v1 SDKでコンパイル済みのパイプラインは、公式の移行ガイド上は「v2バックエンドで変更なしに実行できる」とされています。ただしv1は、2025年3月31日の告知によりバックエンド2.5(2025年4月公開)以降は非推奨です。同年5月21日には「2025年第4四半期の最終リリースでv1を完全に削除し、v1を含む最後のリリースは2026年6月まで延長サポートする」との追記がありました。実際には2026年5月のバックエンド2.16.1と7月の2.17.0で、v1の実行をnamespace単位で遮断する機能が追加されています。つまり2.17系でもv1は動きますが、止める仕組みが整いつつある段階です。v1のコードを使い続ける前提は置かないほうが安全です。
| v1の書き方 | v2での置き換え |
|---|---|
| func_to_container_op/create_component_from_func | @dsl.component |
| dsl.ContainerOp | @dsl.container_component |
| VolumeOp/ResourceOp | kfp-kubernetes(CreatePVC等) |
| from kfp.v2 import dsl | from kfp import dsl |
| 位置引数での呼び出し | キーワード引数が必須 |
| ループ内タスクへの.after() | 実行順指定は.after()を継続、出力集約はdsl.Collected |
| dsl.Condition | dsl.If(Conditionは非推奨) |
| .json拡張子でのコンパイル出力 | .yaml拡張子(YAMLが推奨形式) |
v1のコンポーネントYAMLはcomponents.load_component_from_fileで読み込めますが、v1の軽量コンポーネント用の型(InputTextFile・OutputBinaryFile等)はなくなりました。移行で止まりやすいのは次の3つです。
- 型の緩さに頼った箇所:文字列で数値を渡していた引数はコンパイルエラーになります。
- VolumeOpでの受け渡し:まずArtifactに置き換えられないかを検討し、無理な場合にだけkfp-kubernetesのPVCを使います。
- Vertex AI専用だったv1の機能:
AIPlatformClientやrun_as_aiplatform_custom_jobはなくなりました。投入はgoogle-cloud-aiplatformのPipelineJobで行います。
Agent Platform Pipelines(旧Vertex AI Pipelines)への投入手順と料金
Google Cloudは2026年4月にVertex AIをGemini Enterprise Agent Platformへ改称し、Vertex AI PipelinesもGemini Enterprise Agent Platform Pipelinesになりました。KFP SDK 2.x(kfp>=2,<3)でコンパイルしたIR YAMLを、google-cloud-aiplatformのPipelineJobで投入する手順は改称前と同じです。
from google.cloud import aiplatform
aiplatform.init(project="my-project", location="asia-northeast1")
job = aiplatform.PipelineJob(
display_name="kfp-v2-demo",
template_path="demo.yaml", # KFP SDK 2.x でコンパイルしたIR YAML
pipeline_root="gs://my-bucket/pipeline-root",
parameter_values={"n": 100},
enable_caching=None, # コンパイル済みのタスク別キャッシュ設定を維持
)
job.submit(service_account="[email protected]")
「gcloud ai pipelines run」で検索されることがありますが、gcloudのgcloud aiコマンド群にはGA・beta・alphaのいずれにもpipelinesのサブグループがありません(2026-10-03時点のリファレンスで確認)。CLIから動かしたい場合は、上のPythonスクリプトを実行するか、REST APIでpipelineJobsを作成します。
料金は、パイプライン実行1回あたり0.03米ドルの実行料と、各タスクが使うCompute Engineなどの料金の合計です。セルフホストのKFPのように実行基盤のクラスタを常時維持する必要はありません。ただし、実行していない間もCloud Storageに保存したデータやモデルなどの保管料金は発生します。Agent Platform全体の構成はVertex AI Agent Builder(Gemini Enterprise Agent Platform)の全体像で整理しています。
KFP v2を選ばないほうがよい場面と代わりの選択肢
KFP v2は「コンテナ単位で再現できる学習パイプラインを、Kubernetes上で何本も回す」組織に向いています。次の条件に当てはまるなら、KFPを入れても運用負荷のほうが勝ちます。
- Kubernetesを運用する体制がない:セルフホストのKFPは、Argo Workflows・MySQL・オブジェクトストア・認証を自分で保守する必要があります。Google Cloudを使っているなら、最初からAgent Platform Pipelinesにします。
- 扱うのがデータの移動と集計だけ:ETLや定期バッチが中心なら、Apache AirflowやDagsterのほうが運用の知見とコネクタが揃っています。
- 欲しいのは実験の記録だけ:パラメータと指標の比較が目的なら、パイプライン基盤より実験管理のツールを先に入れるほうが早く効果が出ます。
- パイプラインが1本で月数回しか回らない:実行頻度が低く、承認履歴やデータ・モデルの系譜管理も不要なら、スクリプトとジョブスケジューラーを優先します。
逆に、学習・評価・登録を複数チームが同じ部品で組み立て、どのデータとコードでどのモデルができたかを後から追える状態が必要なら、KFP v2は有力です。ほかのMLOpsツールとの組み合わせ方はMLOpsツール比較を参照してください。
よくある質問
KFPとは何ですか?
Kubeflow Pipelinesの略で、機械学習のワークフローをコンテナのDAGとして定義し、Kubernetes上で実行・記録する仕組みです。PythonのSDK(PyPIのkfp)と、実行を受け持つバックエンド・UIで構成されます。
KFPのv1とv2の違いは何ですか?
v2はデコレーター(@dsl.component・@dsl.pipeline)で書き、Argo Workflowsに依存しないIR YAMLへコンパイルします。ContainerOpやVolumeOpはSDKから外れ、Kubernetes固有の設定はkfp-kubernetesに移りました。v1はバックエンド2.5から非推奨です。
kfp-kubernetesとは何ですか?
KFP v2のタスクに、Secret・PVC・ノードセレクター・トレラレーションなどKubernetes固有の設定を付けるための拡張ライブラリです。最新版はSDKと同じ2.17.0です。公式ドキュメントが対応先として挙げているのはOSSのKFPバックエンドです。
gcloud ai pipelines runは使えますか?
使えません。2026年10月時点のgcloudには、gcloud ai配下にpipelinesのコマンドがありません。google-cloud-aiplatformのPipelineJob(...).submit()か、REST APIで投入します。
KFPのパイプラインはローカルで動かせますか?
動かせます。kfp localのSubprocessRunnerかDockerRunnerをlocal.init()で指定し、パイプライン関数を呼び出します。SDK 2.16.0以降は条件分岐とループも動きますが、キャッシュ・リトライ・リソース指定は手元では効きません。