セクション5 応用 — ML パイプラインの自動化とオーケストレーション
🔧 対象:中堅エンジニア。再訓練ポリシー、検証ステップの実装、Ray の活用、CI/CD/CT 設計の判断軸を押さえる。
1. 再訓練ポリシーの設計
1.1 再訓練のトリガー
| トリガー | 説明 | 適するケース |
|---|---|---|
| スケジュール | 日次/週次/月次の定期 | データ量が安定、変化緩やか |
| データ量閾値 | 新規データが N 件以上で起動 | 新データ流入が不規則 |
| データドリフト検知 | 入力分布変化を Model Monitoring が検知 | 環境変化が起きやすい |
| コンセプトドリフト検知 | 入力→出力の関係が変化 | 行動パターンの変化が起きやすい |
| モデル精度劣化 | 評価メトリクスが閾値割れ | 重要 KPI を直接モニタ |
| ビジネスイベント駆動 | 季節要因・キャンペーン開始 | 季節性・突発要因 |
| オンデマンド | 手動 | アドホック・緊急修正 |
1.2 ハイブリッド戦略
- スケジュール + ドリフト検知:定期的に走らせつつ、ドリフトで臨時実行
- 段階的なトリガー:軽い劣化なら通知、深刻なら自動再訓練
1.3 再訓練の落とし穴
- ループ的な劣化:再訓練データに偏りがあるとさらに歪む
- データ漏洩:将来データが訓練に混入
- 過剰再訓練:データ量が少ないのに頻繁すぎる
- コスト爆発:トリガー条件が緩いと毎日走る
2. データ・モデル検証ステップの実装
2.1 TFDV を使ったデータ検証パターン
import tensorflow_data_validation as tfdv
# 1. ベースラインのスキーマと統計を生成
stats = tfdv.generate_statistics_from_csv("baseline.csv")
schema = tfdv.infer_schema(stats)
# 2. 新規データの統計を計算
new_stats = tfdv.generate_statistics_from_csv("new_data.csv")
# 3. ベースラインと比較してアノマリーを検出
anomalies = tfdv.validate_statistics(new_stats, schema)
tfdv.display_anomalies(anomalies)
2.2 検証で見るもの
- カラム名・型の一致
- 欠損率の急変
- カーディナリティ変化
- 値の分布(平均・分散・パーセンタイル)
- 新規カテゴリの出現
2.3 パイプライン内での配置
[Ingestion] → [TFDV: validate-data]
↓ (pass)
[Preprocess]
↓
[Train]
↓
[TFMA: evaluate]
↓ (合格判定)
[Register & Deploy]
↓ (失敗 / drift) → alert + 再訓練トリガー
2.4 TFMA / Agent Platform Model Evaluation
- スライス別評価(地域 / 年齢層 / etc.)
- 公平性メトリクス(False Positive Rate の属性別差など)
- ベースラインモデルとの比較
3. 訓練/サービング前処理の一貫性
3.1 Training-Serving Skew が起きる原因
- 訓練時はノートブックで pandas、推論時は本番コードで Python から書き直し
- 訓練時は SQL で集計、推論時は API ベース
- データセットの定義 / 欠損補完方法 / one-hot のカテゴリ集合 がズレる
3.2 対策3つ
| アプローチ | 仕組み | 強み |
|---|---|---|
| TFT (TensorFlow Transform) | 訓練時に "Transform graph" を出力、推論時に再利用 | TF エコシステムで完結 |
| Feature Store | 特徴量を中央保管、訓練/推論で同じ取得 | フレームワーク非依存 |
| 共通前処理ライブラリ | Python パッケージとして共有 | 軽量だがバージョン管理が必要 |
3.3 TFT の典型コード
import tensorflow_transform as tft
def preprocessing_fn(inputs):
outputs = {}
outputs["age_normalized"] = tft.scale_to_z_score(inputs["age"])
outputs["category_id"] = tft.compute_and_apply_vocabulary(inputs["category"])
return outputs
出力された Transform graph は TFX BulkInferrer や Agent Platform Inference の前処理層 として配布される。
4. Ray on Agent Platform の応用
4.1 使うべきシーン
- 強化学習 (RLlib):ロボット制御、ゲーム AI、推薦システムのオンライン学習
- 大規模ハイパラチューニング (Ray Tune):数百〜数千試行の並列
- 大規模 LLM ファインチューニング:DeepSpeed / FSDP と統合
- オンライン分散特徴量計算
4.2 構成
[Ray Cluster on Agent Platform]
├─ Head node
├─ Worker nodes (GPU/CPU)
└─ Ray Serve (推論サービング併設可)
4.3 Pipelines / Airflow との関係
- Pipelines / Airflow から Ray ジョブを呼ぶ のが典型
- Ray 自体は タスク粒度の細かい分散実行 に強い
- DAG 全体のオーケストレーションは Pipelines/Airflow が担当
5. CI/CD/CT パイプラインの設計
5.1 全体像
[ Source Repo (Git) ]
│ (push to main / feature branch)
▼
[ Cloud Build trigger ]
│
├─ ① CI: テスト
│ - Unit test
│ - データバリデーション(小規模)
│ - 静的解析
│
├─ ② Container Build
│ - 訓練コードのコンテナ化
│ - Artifact Registry へ push
│
├─ ③ Pipeline Compile & Submit
│ - KFP コンパイル
│ - Agent Platform Pipelines に submit
│
└─ ④ CT: 訓練 → 評価 → モデル検証 → Registry 登録
│
└─ 合格 → CD: 自動デプロイ(カナリア)
└─ 失敗 → アラート / 旧モデルキープ
5.2 CI のチェック項目
- 訓練コードの unit test
- 前処理関数の単体テスト
- データバリデーション(少量サンプル)
- セキュリティスキャン(依存ライブラリ脆弱性)
- Lint / 型チェック
5.3 CT トリガーのパターン
- 新規コードプッシュ → 開発環境で訓練
- 新規データ → ステージング/本番で訓練
- ドリフト検知 → 自動再訓練
- スケジュール → 定期実行
5.4 Cloud Build の連携サービス
- Source Repositories / GitHub / GitLab
- Artifact Registry:コンテナ・パッケージ管理
- Secret Manager:API キー・認証情報
- Cloud Deploy:環境(dev/staging/prod)のプロモーション
6. パイプラインの再実行と部分実行
6.1 Caching
- KFP は ステップごとに入力ハッシュでキャッシュ
- 同じ入力なら再実行不要、開発時に高速化
- 本番では適切に無効化(最新データを使うため)
6.2 失敗時のリトライ
- KFP コンポーネントに
retry_policyを設定 - 一過性の障害(ネットワーク・API quota)に対応
6.3 部分再実行
- パイプラインを途中ステップから再走させたい時は 対象ステップ + 下流のみ submit
- KFP の v2 SDK は柔軟(個別コンポーネントの再実行)
7. パイプラインのスケジュールとイベント駆動
7.1 スケジュール
- Cloud Scheduler → Agent Platform Pipelines
- Pipeline 自体のスケジュール機能(API/Console から)
7.2 イベント駆動
- Pub/Sub → Cloud Function → Pipelines submit
- GCS object create → Eventarc → Pipelines submit
- 新規データ流入で自動実行
7.3 BigQuery 連携
- Scheduled Query → BigQuery → Pub/Sub → Pipelines
- 新規データが入った合図でパイプライン起動
8. 試験での頻出ひっかけ
| シナリオ | 不正解になりがち | 正解の方向 |
|---|---|---|
| 「ML パイプラインを Kubernetes クラスタ管理なしで」 | GKE 自前構築 | Agent Platform Pipelines (サーバーレス) |
| 「複雑な ETL + ML を1つのオーケストレータで」 | Pipelines | Managed Service for Apache Airflow |
| 「強化学習で並列ロールアウト」 | Pipelines / Airflow | Ray on Agent Platform |
| 「軽量な API オーケストレーション」 | Pipelines | Cloud Workflows |
| 「訓練/推論の前処理がずれて精度低下」 | コードを手動同期 | TFT or Feature Store |
| 「データ品質を自動検証してから訓練」 | スキップ | TFDV を組み込む |
| 「モデル評価が前バージョン以下なら本番化阻止」 | デプロイ強行 | モデル検証ステップで条件分岐 |
| 「データドリフトに応じた自動再訓練」 | スケジュール再訓練のみ | Model Monitoring → Pub/Sub → Pipelines (CT) |
| 「コード変更で自動的にパイプライン更新+再訓練」 | 手動 submit | Cloud Build + Pipelines (CI/CD/CT) |
| 「Cloud Composer の新名称は?」 | そのまま | Managed Service for Apache Airflow |
| 「Pipelines に毎回大量パラメータを渡す」 | コードベタ書き | Pipeline parameters + parameterized templates |
| 「失敗ステップだけ再実行したい」 | パイプライン全体再走 | KFP v2 の部分再実行 / キャッシュ |
| 「数千試行の HPT」 | Vizier のみ | Ray Tune (Ray on Agent Platform) |
9. 設計パターン例
9.1 標準的な MLOps レベル 2 構成
[Git: train.py, pipeline.py]
↓ push
[Cloud Build]
├─ unit test
├─ container build → Artifact Registry
├─ pipeline compile
└─ pipeline submit → [Agent Platform Pipelines]
├─ ingestion (BQ)
├─ tfdv-validate
├─ tft-transform
├─ train (Custom Training)
├─ tfma-evaluate
├─ model-validate (vs baseline)
├─ register (Model Registry)
└─ deploy (canary 5%)
↓ monitor
[Model Monitoring]
↓ drift → Pub/Sub
[Re-trigger pipeline (CT)]
9.2 ETL + ML 統合パターン(Airflow)
[Airflow DAG]
├─ BigQuery extract task
├─ Dataflow transform task
├─ Pub/Sub publish task
├─ (parallel) Agent Platform Pipeline trigger task
└─ Slack alert task
9.3 RL 用 Ray パターン
[Pipelines]
├─ data prep (BigQuery)
└─ Ray RLlib job (rollouts + training)
↓
[Model Registry]
↓
[Endpoint deploy]
次は 03_要点と暗記.md で意思決定ツリー・対比表・暗記カードを集めます。