セクション2 応用:データの取り込みと処理 🔧🎯
このファイルは 🔧 実践(中堅) と 🎯 発展(シニア) レベル。 「どう設計判断するか」「トレードオフは何か」「試験のひっかけ」に焦点を当てます。基礎概念は 01_基礎.md を参照。 最重要セクションのため、判断軸とひっかけを特に厚く扱います。
2.1 パイプライン計画の設計判断
🔧 ソース・シンク・処理方式の決定木
質問1: データは「蓄積済み」か「流れ続ける」か?
蓄積済み(有界) → バッチ
流れ続ける(無界) → ストリーミング
質問2: リアルタイム性は必要か?
秒〜分以内の鮮度が必要 → ストリーミング(Pub/Sub + Dataflow)
時間〜日次でよい → バッチ(スケジュール実行)
質問3: 既存のコード資産は?
Spark/Hadoopのコードがある → Dataproc
新規・サーバーレス志向 → Dataflow
コードを書きたくない → Cloud Data Fusion
SQLで完結できる → BigQuery (ELT) / Dataform
🔧 ネットワーキングのベストプラクティス
- Dataflow ワーカーは 外部IPを無効化(
--no_use_public_ips) し、限定公開のGoogleアクセスで Google API へ到達させる(攻撃面の縮小) - サブネットを明示指定し、ファイアウォールでワーカー間通信を許可
- オンプレ連携は Cloud Interconnect / VPN でプライベート接続。公衆インターネット経由は避ける
- ⚠️ 「Dataflowをセキュアに」→ 外部IP無効+限定公開のGoogleアクセス+(必要なら)VPC Service Controls
🔧 暗号化の判断
- 規制要件で鍵を自社管理 → CMEK(Pub/Sub・Dataflow一時データ・BigQuery・GCS すべてに指定可能)
- ⚠️ Dataflow は処理中の一時データにもCMEKを適用できる(パイプライン全体で鍵要件を満たす)
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
- ⚠️ 「リフト&シフトで既存Sparkジョブを動かしたい」→ Dataproc(Dataflowに書き換える必要なし)
- ⚠️ 「サーバーレスでストリーミングとバッチを1つのコードで」→ Dataflow
- ⚠️ 「アナリストがGUIでETLを組みたい」→ Cloud Data Fusion
🔧🎯 Dataflow Shuffle と Streaming Engine
ワーカーVMの負荷をサービス側にオフロードして、スケーラビリティと安定性を高める機能。
| 機能 | 対象 | 効果 |
|---|---|---|
| Dataflow Shuffle | バッチ | シャッフル(GroupBy/Join)をサービス側で実行。VMのCPU/メモリ負荷を軽減、起動高速化 |
| Streaming Engine | ストリーミング | ウィンドウ状態の保持をサービス側へ。ワーカーを軽量化、オートスケールを高速・滑らかに |
- どちらも状態をワーカーから切り離すことで、ワーカーの障害・スケール変更に強くなる
- ⚠️ 「ストリーミングのオートスケールやワーカー安定性を改善」→ Streaming Engine
- ⚠️ 「バッチのシャッフルがVMを圧迫」→ Dataflow Shuffle
- 関連:FlexRS(バッチを低コストで実行する遅延スケジューリング)、Horizontal/Vertical Autoscaling
🎯 ウィンドウ選択の判断
| 要件 | 選ぶウィンドウ |
|---|---|
| 1分ごと・1時間ごとなど定期的な集計 | 固定(Fixed/Tumbling) |
| 直近N分の移動平均・傾向 | スライディング(Sliding/Hopping) |
| ユーザーの操作のまとまり・無活動で区切る | セッション(Session) |
- ⚠️ 「30分操作がなければセッション終了とみなす」→ セッションウィンドウ(gap=30分)
- ⚠️ 「5分間の平均を1分ごとに更新」→ スライディング(window=5分, period=1分)
- ⚠️ 「重複なく毎分カウント」→ 固定ウィンドウ
🎯 ウォーターマーク・遅延データの設計
- イベント時刻で集計するのが原則(処理時刻だと遅延・順序逆転で結果が歪む)
- ⚠️ 「ネットワーク遅延でデータが順不同に届くが正しく集計したい」→ イベント時刻+ウォーターマーク
- 遅延データへの対応:
- allowed lateness(許容遅延) を設定 → 遅れて来たデータでウィンドウ結果を再計算・更新
- 早期トリガーで途中経過を出し、確定後に更新する運用も可能
- 許容を超える超遅延は破棄 or デッドレターへ
- ⚠️ 「集計後に遅れて届くデータも反映したい」→ トリガー+許容遅延(accumulating mode)
🔧🎯 Exactly-once の実現
- Pub/Sub(at-least-once、重複あり得る)+ Dataflow(重複排除・チェックポイント)→ Exactly-once
- ⚠️ 「重複なくちょうど1回処理したい」→ Dataflow の Exactly-once 処理(Pub/Subの重複をDataflowが吸収)
- BigQuery への書き込みは Storage Write API で 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%程度) |
| 運用 | 完全フルマネージド・無設定 | キャパシティ計画が必要 |
| 選ぶ条件 | 一般用途・可用性重視・運用最小 | コスト最優先・高スループットで予測可能 |
- ⚠️ 「とにかくコストを抑えたい、スループットは予測でき、ゾーン単位でよい」→ Pub/Sub Lite
- ⚠️ 「グローバルで高可用・運用したくない」→ Pub/Sub(既定はこちら)
- 補足:Pub/Sub Lite は将来的な扱いに注意が必要だが、試験上は「低コスト・要キャパシティ計画」の対比で問われる
🎯 Kafka 資産の扱い
- 既存Kafkaがある → Managed Service for Apache Kafka や、Dataflow の Kafka コネクタ(KafkaIO)で連携
- ⚠️ 「オンプレKafkaの資産を活かしつつGCPで運用」→ マネージドKafka or Pub/Sub への移行を検討(要件次第)
🎯 AI エンリッチメントの設計
- パイプライン内で外部モデルを呼ぶとレイテンシ・コストが増える → バッチで非同期化、結果をキャッシュ
- ⚠️ 「BigQueryのテキスト列を要約/分類したい(SQLで)」→ BigQuery ML の
ML.GENERATE_TEXT - ⚠️ 「埋め込みを生成してベクトル検索」→
ML.GENERATE_EMBEDDING+ Vector Search
🔧 データ取得・形式の判断
- ⚠️ 大量・型安全・高速ロード → Avro / Parquet(自己記述・圧縮・列指向)。CSVはスキーマ・型に弱い
- ⚠️ 「ロードせずGCS上のファイルを直接分析」→ 外部テーブル / BigLake
- ⚠️ 「定型のGCS→BQ取り込みを素早く」→ Dataflow テンプレート(Google提供)
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
- ⚠️ 「複雑な依存関係のあるDAGを管理」→ Cloud Composer
- ⚠️ 「数個のAPIを順番に呼ぶ軽量な連携を低コストで」→ Workflows
- ⚠️ 「ジョブを毎朝6時に起動するだけ」→ Cloud Scheduler(Composerを使うほどではない)
🔧 Composer 運用のベストプラクティス
- DAGコードは Git で管理し、Cloud Build で GCS の
dags/フォルダへ自動デプロイ - 重い処理はDAG内で直接実行せず、Dataflow/Dataproc/BigQueryへオフロード(Workerを軽く保つ)
- 冪等性(idempotency)を持たせ、リトライしても二重処理にならない設計に
- ⚠️ 「Composerのワーカーが重い処理で詰まる」→ 処理を外部サービスへオフロードし、Composerはオーケストレーションに専念
🔧🎯 CI/CD の設計(Cloud Build)
- トリガー:Gitへの push / PR で
cloudbuild.yamlを実行 - 段階デプロイ:dev → staging → prod。各段でテストを通過したものだけ昇格
- Dataflow:Flex テンプレートをビルドし Artifact Registry に登録 → 環境ごとにパラメータ実行
- インフラ:Terraform を Cloud Build から適用(IaC)
- ⚠️ 「パイプラインのデプロイを自動化・再現可能に」→ Cloud Build による CI/CD + テンプレート/IaC
🎯 ひっかけ・判断ポイント総まとめ
- ⚠️ Dataflow vs Dataproc:「既存Spark/Hadoop」が出たら Dataproc、「サーバーレス/Beam/ストリーミング」なら Dataflow
- ⚠️ Composer vs Workflows:「複雑なDAG/データジョブ」は Composer、「軽量API連携/低コスト」は Workflows
- ⚠️ Pub/Sub vs Pub/Sub Lite:「コスト最優先+キャパシティ予測可」は Lite、既定は Pub/Sub
- ⚠️ ウィンドウ:定期集計=固定、移動平均=スライディング、無活動で区切る=セッション
- ⚠️ イベント時刻 vs 処理時刻:順序乱れ・遅延を正しく扱うならイベント時刻+ウォーターマーク
- ⚠️ Shuffle vs Streaming Engine:バッチのシャッフル負荷=Shuffle、ストリーミングの安定/オートスケール=Streaming Engine
- ⚠️ Exactly-once:Pub/Subは at-least-once(重複あり)、Dataflowが重複排除して Exactly-once を実現
- ⚠️ GUI/ノーコードが出たら Cloud Data Fusion
- ⚠️ デッドレター=処理失敗メッセージの隔離、シーク=過去への巻き戻し再処理
🎯 統合シナリオ演習(考え方の練習)
シナリオ:あるEC企業。①Webとモバイルから毎秒数万件のクリック/購入イベントが発生。②リアルタイムで「直近5分の売上トレンド」をダッシュボードに表示したい。③ユーザーの操作セッション単位の分析もしたい。④イベントは順不同・遅延ありで届く。⑤一部メッセージが壊れていてもパイプラインを止めたくない。⑥同じユーザーのイベントは順序を保ちたい。⑦夜間に当日分をBigQueryで集計しレポートを生成、複数ジョブの依存がある。⑧パイプラインは自動デプロイしたい。
設計の骨子(解答例)
- 取り込み:イベントは Pub/Sub で受ける(急な流入をバッファ、疎結合)。同一ユーザーの順序保持に 順序指定キー(Ordering Key) を使用
- ストリーミング処理:Dataflow(Apache Beam, ストリーミング) で処理。Streaming Engine を有効化してオートスケールと安定性を確保
- 直近5分のトレンド:スライディングウィンドウ(window=5分, period=1分) で集計
- セッション分析:セッションウィンドウ(無活動gapで区切る) を別途適用
- 順不同・遅延:イベント時刻+ウォーターマークで集計し、許容遅延(allowed lateness) とトリガーで遅延データを反映
- 壊れたメッセージ:デッドレタートピック/テーブルへ隔離し、本処理は継続(Exactly-once は Dataflow が担保)
- シンク:集計結果を BigQuery へ(Storage Write API で Exactly-once)。ダッシュボードは Looker/Looker Studio(+必要なら BI Engine)
- 夜間バッチ+依存管理:Cloud Composer(Airflow DAG) で「集計→検証→レポート→通知」を依存関係付きでオーケストレーション
- CI/CD:DAGとDataflowテンプレートを Git 管理し、Cloud Build で dev→staging→prod へ自動デプロイ
この「リアルタイム要件=Pub/Sub+Dataflow+ウィンドウ」「複雑な依存のバッチ=Composer」「自動化=Cloud Build」の組み合わせが、本セクションの設計問題の典型パターンです。
まとめ:このセクションの設計判断の型
- 処理方式は バッチ/ストリーミング → サービスは Dataflow/Dataproc/Data Fusion/BigQuery を要件で選ぶ
- ストリーミングは ウィンドウ・ウォーターマーク・トリガー・遅延データ・Exactly-once が判断の核
- スケール/安定性は Streaming Engine(ストリーミング)/ Shuffle(バッチ)
- 取り込みの緩衝は Pub/Sub(順序指定キー・デッドレター・シーク/コストなら Lite)
- オーケストレーションは 複雑=Composer / 軽量=Workflows、デプロイは Cloud Build
→ 03_要点と暗記.md で記憶を固めましょう。