セクション5 基礎 — ML パイプラインの自動化とオーケストレーション
📘 対象:新人エンジニア。E2E ML パイプラインの全体像、主要なオーケストレータ、MLOps 成熟度の3レベルを押さえる。
1. なぜパイプラインが必要か
手動の ML 開発の問題:
- 「ノートブックで学習 → 結果を Slack で共有」だと再現できない
- 新しいデータが来たときに 同じ前処理・同じ学習設定 で動かすのが大変
- 本番モデルがどの訓練データ・どのコード・どのハイパラから生まれたか追跡不能
- ロールバック時に旧モデル成果物が見つからない
パイプライン化のメリット:
- 再現性・系統管理 (lineage)
- 自動化(夜間/イベント駆動)
- 部分再実行(失敗ステップだけ再走)
- 監査・ガバナンス対応
2. ML パイプラインの典型構成
[1. データ取り込み]
↓
[2. データ検証](スキーマ/統計/異常)
↓
[3. 前処理・特徴量化](TFT / Feature Store)
↓
[4. 訓練]
↓
[5. 評価]
↓
[6. モデル検証](性能 / バイアス / SLO 合格判定)
↓
[7. Model Registry に登録]
↓
[8. 推論基盤にデプロイ](カナリア → 本番)
↓
[9. モニタリング](ドリフト / 精度劣化)
↓ (劣化検知)
[再訓練トリガー]
各ステップが コンポーネント(再利用可能なユニット)として独立しているのが理想。
3. Agent Platform Pipelines(旧 Vertex AI Pipelines)
3.1 概要
- Kubeflow Pipelines (KFP) または TFX で書いた DAG を サーバーレス で実行
- 各ステップ = コンテナ。入出力アーティファクトを ML Metadata に自動記録
- スケジュール実行・イベント駆動(Pub/Sub / Eventarc)
3.2 KFP v2 基本構造(Python SDK)
from kfp import dsl, compiler
@dsl.component(packages_to_install=["pandas", "scikit-learn"])
def train(data_path: str, model_path: dsl.OutputPath("Model")):
import pandas as pd, joblib
from sklearn.linear_model import LogisticRegression
df = pd.read_csv(data_path)
model = LogisticRegression().fit(df.drop("y", axis=1), df["y"])
joblib.dump(model, model_path)
@dsl.pipeline(name="churn-pipeline")
def pipeline(data_path: str):
train(data_path=data_path)
compiler.Compiler().compile(pipeline, "pipeline.json")
3.3 強み
- サーバーレス(クラスタ管理不要)
- ML Metadata と統合(系統が自動記録)
- Google Cloud サービスとの統合コンポーネントが豊富
- KFP 標準準拠、OSS との互換性高い
3.4 いつ使う
- ML パイプラインの 第一候補
- DAG がシンプル〜中程度
- Kubernetes 詳細を意識したくない
4. Managed Service for Apache Airflow(旧 Cloud Composer)
4.1 概要
- Apache Airflow をマネージド提供
- 旧名は Cloud Composer。新ガイドでは "Managed Service for Apache Airflow" 表記
- DAG を Python で記述、豊富な Operator
4.2 強み
- 複雑な依存関係に強い(条件分岐・ループ・タスクグループ)
- BigQuery / GCS / Dataflow / GKE 等の 大量の Operator
- ETL と ML を統合したパイプラインに向く
- 既存 Airflow 資産(オンプレ含む)を移行可能
4.3 いつ使う
- ETL + ML を1つのオーケストレータで管理
- 既存 Airflow ワークフローを Google Cloud に移行
- 複雑な依存・条件分岐
4.4 Cloud Composer 1 vs Composer 2 vs Composer 3
- Composer 3 が最新、サーバーレス度が上がりコスト効率改善
- DAG コードは互換
5. Ray on Agent Platform(新ガイドで追加)🆕
5.1 Ray とは
- 分散 Python フレームワーク
- 強化学習(RLlib)、ハイパラチューニング(Tune)、データ並列(Train)、サービング(Serve)
- 単一クラスタで多様な分散 ML ワークロード
5.2 Ray on Agent Platform
- マネージド Ray クラスタを Agent Platform 上で起動
- 強化学習、大規模ハイパラチューニング、大規模分散訓練、リアルタイム特徴量計算で使う
- 通常の ML パイプラインよりも柔軟な分散制御が必要なケース
5.3 いつ使う
- 強化学習 (RL) のロールアウト
- 大規模 分散ハイパラチューニング (Ray Tune)
- LLM の 分散ファインチューニング(DeepSpeed/FSDP と統合)
- Pipelines / Airflow では表現が難しい複雑な分散ワークロード
6. Cloud Workflows
6.1 概要
- YAML ベースのサーバーレスオーケストレーション
- HTTP API・Google Cloud サービスを順次/並列呼び出し
- 軽量・低コスト・低レイテンシ起動
6.2 いつ使う
- 軽量な API オーケストレーション(数ステップの呼び出し)
- パイプラインのトリガーやデプロイの自動化
- 「Pub/Sub → Cloud Function → BigQuery → Slack 通知」のような短いフロー
6.3 Pipelines / Airflow との違い
- Workflows は タスクの単位がジョブ(実行は別サービスで)
- Pipelines / Airflow は DAG 全体を自分で管理
- Workflows = orchestration glue、Pipelines/Airflow = workflow engine
7. オーケストレータ早見表
| Agent Platform Pipelines | Managed Airflow | Ray on Agent Platform | Cloud Workflows | |
|---|---|---|---|---|
| 主用途 | ML パイプライン | ETL+ML/汎用ワークフロー | 分散 ML / RL | 軽量 API 連携 |
| 言語 | Python (KFP v2) | Python (Airflow DSL) | Python (Ray API) | YAML |
| 強み | サーバーレス・系統管理 | Operator 豊富・柔軟 | 分散制御 | 軽量・低コスト |
| 第一候補 | ML 専用パイプライン | ETL+ML、既存資産 | RL/HPT/分散ML | API オーケストレーション |
8. データ・モデル検証ステップ
8.1 データ検証
- スキーマ検証:カラム名・型・nullable のチェック
- 統計検証:平均・分散・カーディナリティの想定範囲
- 異常検出:欠損率急増、新規カテゴリ出現
ツール:
- TFDV (TensorFlow Data Validation):スキーマ + 統計ベースライン
- Great Expectations:DSL で期待値定義
- BigQuery ASSERT 文 / dbt tests:SQL ベース
8.2 モデル検証
- オフライン評価:ホールドアウトデータでメトリクス確認
- 公平性 (fairness) 検証:人種・性別など属性別の精度差
- SLO 合格判定:レイテンシ、精度の下限を満たすか
- 退行検証:前バージョンを下回らないか
- 必須:本番デプロイ前のゲート条件
ツール:
- TFMA (TensorFlow Model Analysis)
- Agent Platform Model Evaluation
9. CI/CD/CT の概念
9.1 用語
- CI (Continuous Integration):コード/データの統合と自動テスト
- CD (Continuous Delivery / Deployment):本番への自動配信
- CT (Continuous Training):データ変化に応じた自動再訓練
9.2 MLOps 成熟度モデル(Google 公式)
| レベル | 名称 | 特徴 |
|---|---|---|
| Level 0 | 手動 | 手動でノートブック実行、人手でデプロイ |
| Level 1 | ML パイプライン自動化 | パイプラインで自動訓練・自動デプロイ |
| Level 2 | CI/CD パイプライン自動化 | コードの変更で パイプライン自体 も自動更新 + CT |
試験頻出: 「レベル X はどう違うか」
9.3 Cloud Build の役割
- コード変更 (git push) をトリガーに テスト + ビルド + デプロイ を自動実行
- 訓練コードや前処理コードのテスト → コンテナビルド → Artifact Registry プッシュ → パイプライン定義更新
9.4 典型的な CI/CD/CT 構成
[Developer commits code]
↓ (Cloud Build trigger)
[Unit tests + Lint]
↓
[Container build & push (Artifact Registry)]
↓
[Pipeline definition compile (KFP)]
↓
[Pipeline submit (Agent Platform Pipelines)]
↓
[Training → Evaluation → Model Registry → Canary deploy]
↓
[Monitoring (drift detection)]
↓ (drift detected)
[Auto-trigger pipeline → CT]
10. このセクションで覚えるキーワード
- Agent Platform Pipelines (KFP v2):サーバーレス ML パイプライン
- Managed Service for Apache Airflow (旧 Cloud Composer):汎用ワークフロー
- 🆕 Ray on Agent Platform:分散 ML / RL / HPT
- Cloud Workflows:軽量 API オーケストレーション
- TFX / TFDV / TFMA / TFT:データ・モデル検証・前処理
- Cloud Build / Artifact Registry:CI/CD パイプライン
- MLOps Level 0 / 1 / 2:成熟度モデル
- CT (Continuous Training):データ駆動の自動再訓練
次は 02_応用.md で再訓練ポリシーと検証ステップの応用、Ray の詳細を扱います。