なぜデータパイプラインは成長とともに破綻するのか

データパイプラインの多くは、単一のスクリプトから始まります。しかしデータセットが増え、ワークフローが3本、5本、10本と増えた瞬間から問題が顕在化します。

  • 変換ロジックのコピペ地獄: format_date 関数を1つ修正するために、7本のスクリプトを同時に修正する必要があります。
  • 意図がコードに埋もれる: 「このパイプラインは結局何をしているのか」を知るには、スクリプトを最初から最後まで読むしかありません。
  • 検証が遅すぎる: スキーマエラー、権限問題、誤ったカラムマッピングが実行中の途中で発覚します。医療・金融・ライフサイエンスのような規制環境では、これはそのまま監査リスクに直結します。

根本原因はただ一つです。「何を(What)」と「どのように(How)」が同じスクリプト内に混在していることです。本記事ではこの2つを分離する**Specification-Driven Composition(仕様駆動コンポジション)**パターンを扱います。根拠資料は AWS Architecture Blog 原文で確認できます。

Architecture diagram showing three-layer specification-driven composition pattern for data pipelines on AWS cloud Programming Illustration

3層アーキテクチャ: Intent / Composition / Processing

仕様駆動コンポジションは、ワークフローを3つのレイヤーに分割します。

レイヤー役割成果物
Intent(意図)ワークフローが「何を」すべきかを定義JSON/YAML Specification
Composition(組成)仕様を検証し、実行可能なパイプラインへコンパイルAmazon States Language (ASL)
Processing(実行)再利用可能な変換ステップを順次実行Lambda Capability Processors

中核となる4つのコンポーネントは以下の通りです。

  1. Specification — データセット、マッピング、変換を記述した宣言的ドキュメント(バージョン管理対象)
  2. Composer — 仕様を読み、capabilityの存在を検証した上でASLへコンパイル(変換自体は行わない)
  3. Capability Registry — 再利用可能な変換関数のメタデータストア(OpenSearchベース、全文検索・セマンティック検索対応)
  4. Capability Pipeline — 組成されたStep Functionsが各Lambdaプロセッサを順次呼び出し

仕様の例(JSON)

{
  "source": ["raw_orders"],
  "target": ["curated_orders"],
  "mappings": [
    {
      "source_field": "order_date",
      "target_field": "order_date_iso",
      "capability": "format_date@1.2.0",
      "params": { "format": "ISO8601" }
    },
    {
      "source_field": "amount",
      "target_field": "amount_usd",
      "capability": "normalize_currency@2.0.1",
      "params": { "target_currency": "USD" }
    }
  ],
  "preprocessing": [
    { "capability": "standardize_columns@1.0.0" }
  ]
}

Composerの中核ロジック(擬似コード)

# 仕様を読み込み、capabilityを検証し、Step Functionsを起動します。
def compose_pipeline(spec_s3_key: str) -> str:
    spec = load_json_from_s3(spec_s3_key)
    validate_against_schema(spec)  # スキーマ検証失敗時は即座に中断

    for mapping in spec["mappings"]:
        cap_id = mapping["capability"]  # 例: "format_date@1.2.0"
        meta = opensearch.lookup(cap_id)  # ARN、I/Oスキーマ、権限境界を照会
        if meta is None:
            raise CapabilityNotFound(cap_id)  # ランタイム前に失敗させる

    asl = build_state_machine(spec)  # 再利用可能なASLへコンパイル
    return stepfunctions.start_execution(asl)

重要なのは、Composerは変換を実行しないという点です。担当するのは「組成」のみです。この分離が規制環境で特に強力な理由は、ドメインユーザーが仕様のみを作成し、実行コードはシステムが生成するという**職務分掌(Separation of Duties)**の原則を自然に満たせるからです。

Serverless AWS workflow with Lambda composer, Step Functions orchestration and OpenSearch capability registry Developer Related Image

このパターンが輝く場面 vs. 過剰設計になる場面

適しているケース

  • 規制産業のデータ提出パイプライン: 例えば臨床試験データをFDA提出用のSDTM(Study Data Tabulation Model)形式へ変換するケース。アナリストが仕様のみを記述すれば、Composerが検証済みcapabilityでパイプラインを組成するため、監査対応時間が大幅に短縮されます。
  • マルチソース統合: 異なるソースシステムから月次財務レポートを生成する場合のように、バリアントが3つ以上あるパイプライン。
  • 再利用可能なETLフレームワーク: 変換ロジックを一度実装し、複数のワークフローで再利用したい場合。

過剰設計になるケース

  • 単発の変換: 一度実行して破棄するスクリプトにComposerを付けるのは無駄です。
  • ワークフローが3〜5本未満: 重複排除の効果よりも、Composer・Registryの運用コストの方が大きくなります。

セキュリティ・規制観点のチェックリスト

  • 仕様・データ用S3バケットに**SSE-KMS(カスタマー管理キー)**を適用し、バケットポリシーで aws:SecureTransport を強制
  • OpenSearchドメインで保存時の暗号化+ノード間暗号化を有効化
  • 機密度タグ付け: 仕様に "sensitivity": "PHI" のようにフィールド機密度をタグ付けし、capabilityごとの挙動(保持/除去/生成)をRegistryに宣言 → Composerが自動的にマスキング成果物(Lake Formationカラムグラント等)を生成
  • Capabilityバージョンの固定: 仕様で format_date@1.2.0 のように明示的なバージョン参照 → 再現性を確保

日本の開発エコシステムにおける適用文脈

日本の金融・医療系システムでは、このパターンは特に有効です。金融庁検査やPMDA対応が必要なプロジェクトにおいて、「このパイプラインがどのような変換を行っているか」を仕様書1枚で説明できることは大きな武器になります。ただし国内環境はレガシーオンプレミス+バッチ中心のケースが多く、Step Functionsの代わりにAirflow DAGへコンパイルする変形の方が現実的です。Composerの「ASL生成」部分を「DAG生成」に置き換えれば、概念はそのまま維持できます。

Data analyst reviewing JSON specification mapping source fields to target fields in a declarative pipeline Software Concept Art

まとめ: パイプラインを「設定作業」に変える投資

仕様駆動コンポジションの本質は、エンジニアリング作業を設定作業へ転換することです。再利用可能なcapabilityライブラリと規律ある仕様フォーマットに初期投資を行えば、以降の新規パイプラインはコード修正なしで仕様記述のみで作成できます。

実務で即座に始めるには、以下の順序を推奨します。

  1. 既存パイプラインを1本選び、仕様として文書化 — コードを変えずに仕様だけ書くだけでもインサイトが得られます。
  2. 再利用可能な変換関数を3〜5個capabilityとして登録 — Registryから作り始めてください。
  3. Composerプロトタイプを作成 — S3アップロード → Lambda → Step Functionsの流れで最小実装。
  4. 月次財務レポートのようにバリアントが3つ以上あるパイプラインに適用 — オンボーディング時間を測定し、効果を定量化してください。

このパターンの限界

  • **Composer自体が単一障害点(SPOF)**になり得ます。Composerの障害は、すべてのパイプラインの新規デプロイを麻痺させます。
  • 仕様スキーマの進化管理が新たな負担となります。スキーマのバージョニングとマイグレーション戦略なしでは、仕様がそのまま技術的負債になります。
  • デバッグ難易度の上昇: 「なぜこのフィールドがこう変換されたのか」を追跡するには、仕様 → ASL → Lambdaの3段階を横断する必要があります。可観測性(OpenTelemetryトレーシング)を最初から設計に組み込んでください。

次のステップ学習方向

あわせて読みたい記事

本コンテンツは、信頼性の高い情報源をもとにAIツールを活用して作成され、編集者によるレビューを経て公開されています。専門家によるアドバイスの代替となるものではありません。