02_応用

セクション2 応用:データの取り込みと処理 🔧🎯

このファイルは 🔧 実践(中堅)🎯 発展(シニア) レベル。 「どう設計判断するか」「トレードオフは何か」「試験のひっかけ」に焦点を当てます。基礎概念は 01_基礎.md を参照。 最重要セクションのため、判断軸とひっかけを特に厚く扱います。


2.1 パイプライン計画の設計判断

🔧 ソース・シンク・処理方式の決定木

質問1: データは「蓄積済み」か「流れ続ける」か?
  蓄積済み(有界) → バッチ
  流れ続ける(無界) → ストリーミング

質問2: リアルタイム性は必要か?
  秒〜分以内の鮮度が必要 → ストリーミング(Pub/Sub + Dataflow)
  時間〜日次でよい → バッチ(スケジュール実行)

質問3: 既存のコード資産は?
  Spark/Hadoopのコードがある → Dataproc
  新規・サーバーレス志向 → Dataflow
  コードを書きたくない → Cloud Data Fusion
  SQLで完結できる → BigQuery (ELT) / Dataform

🔧 ネットワーキングのベストプラクティス

🔧 暗号化の判断


2.2 パイプライン構築の設計判断(最重要)

🔧🎯 Dataflow vs Dataproc vs Cloud Data Fusion(頻出の使い分け)

観点 Dataflow Dataproc Cloud Data Fusion
モデル Apache Beam(サーバーレス) Spark/Hadoop(クラスタ) GUI(内部でDataproc/Spark)
運用 クラスタ管理不要・自動スケール クラスタを管理(ephemeral推奨) クラスタ管理不要(裏で起動)
コード Java/Python/Go 既存Spark/PySpark/Hive ノーコード/ローコード
得意 ストリーミング+バッチ統一 既存OSS資産の移行 迅速なETL・非エンジニア
選ぶ条件 新規・サーバーレス・ストリーミング Spark/Hadoopコードがある コードを書かない/書けない
【一発判定】
  「ストリーミング」「サーバーレス」「Apache Beam」 → Dataflow
  「既存のSpark/Hadoop」「Hive」「移行」          → Dataproc
  「GUI」「ノーコード」「ドラッグ&ドロップ」        → Cloud Data Fusion
  「SQLだけで」「BigQuery内のELT」                → BigQuery / Dataform

🔧🎯 Dataflow Shuffle と Streaming Engine

ワーカーVMの負荷をサービス側にオフロードして、スケーラビリティと安定性を高める機能。

機能 対象 効果
Dataflow Shuffle バッチ シャッフル(GroupBy/Join)をサービス側で実行。VMのCPU/メモリ負荷を軽減、起動高速化
Streaming Engine ストリーミング ウィンドウ状態の保持をサービス側へ。ワーカーを軽量化、オートスケールを高速・滑らかに

🎯 ウィンドウ選択の判断

要件 選ぶウィンドウ
1分ごと・1時間ごとなど定期的な集計 固定(Fixed/Tumbling)
直近N分の移動平均・傾向 スライディング(Sliding/Hopping)
ユーザーの操作のまとまり・無活動で区切る セッション(Session)

🎯 ウォーターマーク・遅延データの設計

🔧🎯 Exactly-once の実現

🔧🎯 Pub/Sub の重要機能(頻出)

機能 説明 使いどころ
順序指定キー(Ordering Key) 同じキーのメッセージを順序保証して配信 キーごとに順序が重要なイベント(例: 同一ユーザーの操作順)
デッドレタートピック(DLQ) 規定回数失敗したメッセージを退避 処理不能メッセージの隔離、パイプライン停止防止
シーク(Seek) 過去のタイムスタンプ/スナップショットへ再生位置を戻す 障害後の再処理、バグ修正後の再取り込み
メッセージ保持(Retention) 確認済みでも一定期間メッセージを保持 シークでの再処理を可能にする
フィルタ(Filter) 属性でサブスクリプション配信を絞る 必要なメッセージだけ受信
スキーマ(Schema) Avro/Protobufでメッセージ構造を検証 データ品質の担保

🔧🎯 Pub/Sub vs Pub/Sub Lite

観点 Pub/Sub Pub/Sub Lite
スコープ グローバル(マルチリージョン自動) ゾーン/リージョン単位
キャパシティ 自動プロビジョニング 手動でスループット/ストレージを事前予約
コスト 標準 大幅に安い(〜80%程度)
運用 完全フルマネージド・無設定 キャパシティ計画が必要
選ぶ条件 一般用途・可用性重視・運用最小 コスト最優先・高スループットで予測可能

🎯 Kafka 資産の扱い

🎯 AI エンリッチメントの設計

🔧 データ取得・形式の判断


2.3 デプロイと運用化の設計判断

🔧🎯 Cloud Composer vs Workflows(頻出の使い分け)

観点 Cloud Composer(Airflow) Workflows
本質 データパイプラインのオーケストレータ サーバーレスなAPI/サービス連携
定義 Python(DAG) YAML/JSON
依存関係 複雑な依存・分岐・リトライに強い 順次/簡単な分岐
コスト/運用 常時稼働の環境(コストあり) 実行課金・インフラ管理不要
得意 多数タスク・スケジュール・データ統合 軽量イベント駆動・マイクロサービス連携
選ぶ条件 複雑なETL/ELTの依存管理 軽量・低コスト・APIの連携
【一発判定】
  「DAG」「Airflow」「複雑な依存」「多数のデータジョブ」 → Cloud Composer
  「軽量」「サーバーレス」「API/HTTPの連携」「低コスト」  → Workflows
  「単純な定期トリガー(cron)だけ」                    → Cloud Scheduler

🔧 Composer 運用のベストプラクティス

🔧🎯 CI/CD の設計(Cloud Build)


🎯 ひっかけ・判断ポイント総まとめ


🎯 統合シナリオ演習(考え方の練習)

シナリオ:あるEC企業。①Webとモバイルから毎秒数万件のクリック/購入イベントが発生。②リアルタイムで「直近5分の売上トレンド」をダッシュボードに表示したい。③ユーザーの操作セッション単位の分析もしたい。④イベントは順不同・遅延ありで届く。⑤一部メッセージが壊れていてもパイプラインを止めたくない。⑥同じユーザーのイベントは順序を保ちたい。⑦夜間に当日分をBigQueryで集計しレポートを生成、複数ジョブの依存がある。⑧パイプラインは自動デプロイしたい。

設計の骨子(解答例)

  1. 取り込み:イベントは Pub/Sub で受ける(急な流入をバッファ、疎結合)。同一ユーザーの順序保持に 順序指定キー(Ordering Key) を使用
  2. ストリーミング処理Dataflow(Apache Beam, ストリーミング) で処理。Streaming Engine を有効化してオートスケールと安定性を確保
  3. 直近5分のトレンドスライディングウィンドウ(window=5分, period=1分) で集計
  4. セッション分析セッションウィンドウ(無活動gapで区切る) を別途適用
  5. 順不同・遅延イベント時刻+ウォーターマークで集計し、許容遅延(allowed lateness)トリガーで遅延データを反映
  6. 壊れたメッセージデッドレタートピック/テーブルへ隔離し、本処理は継続(Exactly-once は Dataflow が担保)
  7. シンク:集計結果を BigQuery へ(Storage Write API で Exactly-once)。ダッシュボードは Looker/Looker Studio(+必要なら BI Engine)
  8. 夜間バッチ+依存管理Cloud Composer(Airflow DAG) で「集計→検証→レポート→通知」を依存関係付きでオーケストレーション
  9. CI/CD:DAGとDataflowテンプレートを Git 管理し、Cloud Build で dev→staging→prod へ自動デプロイ

この「リアルタイム要件=Pub/Sub+Dataflow+ウィンドウ」「複雑な依存のバッチ=Composer」「自動化=Cloud Build」の組み合わせが、本セクションの設計問題の典型パターンです。


まとめ:このセクションの設計判断の型

03_要点と暗記.md で記憶を固めましょう。