セクション2:データの取り込みと処理
データパイプラインの中核。「取り込む → 処理する → 運用する」の3工程を、サービスの役割・ストリーミングの仕組み・設計判断・試験のひっかけまで一気通貫で押さえます。
出題 ~25%
★ 最重要セクション
📘 基礎
🔧 応用
🎯 要点と暗記
🔑 TL;DR(このセクションの核)
- パイプラインは ソース→取り込み→変換→シンク + オーケストレーション + CI/CD で考える
- 処理サービスは動詞で覚える:Pub/Sub=受け取る / Dataflow=処理する(サーバーレス) / Dataproc=Sparkを動かす / Data Fusion=GUIで組む / BigQuery=SQLで変換
- ストリーミングの核:ウィンドウ(固定/スライディング/セッション)・ウォーターマーク・トリガー・遅延データ・Exactly-once
- オーケストレーション:複雑な依存=Composer / 軽量連携=Workflows / 単純なcron=Cloud Scheduler、自動デプロイは Cloud Build
🎯 学習目標チェックリスト
📖 学習コンテンツ
📘 基礎レベル
💡 このレベルのゴール
各サブトピックの「概念」と「サービスの役割」を理解すること。設計判断やトレードオフは「応用」タブで扱います。
このセクションは出題比重 ~25% で最重要。サービスの役割と使い分けを確実に押さえましょう。
2.0 全体像:パイプラインの3工程
データパイプラインは「取り込む → 処理する → 運用する」の3工程で考えます。シラバスの2.1/2.2/2.3に対応します。
┌──────────┐ ┌──────────┐ ┌──────────┐ ┌──────────┐
│ ソース │──►│ 取り込み │──►│ 処理/変換 │──►│ シンク │
│ (Source) │ │ (Ingest) │ │(Transform)│ │ (Sink) │
└──────────┘ └──────────┘ └──────────┘ └──────────┘
DB/ログ Pub/Sub Dataflow BigQuery
IoT/SaaS GCS Dataproc GCS
Kafka Storage TS Data Fusion Bigtable
2.1 計画 ──────────────────────────────────────────►
2.2 構築(クレンジング・変換・サービス選定)
2.3 運用(オーケストレーション=Composer/Workflows、CI/CD=Cloud Build)
- 2.1 計画:どこから(ソース)・どこへ(シンク)・どう変換・どう束ねる(オーケストレーション)かを設計
- 2.2 構築:実際にパイプラインを作る。サービス選定と変換ロジックが核心
- 2.3 運用化:定期実行・依存管理・自動デプロイ(CI/CD)
2.1 データパイプラインの計画
① ソース(Source)とシンク(Sink)
パイプラインの入口と出口を定義します。
【ソース=データの出どころ】 【シンク=データの行き先】
- DB(Cloud SQL, オンプレRDB) - BigQuery(分析)
- イベント/ログ(アプリ, IoT) - Cloud Storage(データレイク)
- メッセージ(Pub/Sub, Kafka) - Bigtable(低レイテンシ参照)
- ファイル(GCS, S3, オンプレ) - Pub/Sub(後段への再配信)
- SaaS(Salesforce 等) - Cloud SQL / Spanner(OLTP書き戻し)
- 「何が入って何が出るか」を最初に決めると、必要なサービスが見えてくる
- 同じソース・シンクでも、バッチかストリーミングかで構成が変わる(→ 2.2で詳説)
② 変換(Transformation)とオーケストレーションのロジック
| ロジック | 意味 | 例 |
| 変換 | データの形を変える処理 | 整形・フィルタ・集計・結合・型変換・エンリッチメント |
| オーケストレーション | 複数ジョブの実行順序・依存・スケジュールの制御 | 「抽出→変換→ロード→通知」を順番に・依存を守って実行 |
- 変換は Dataflow / Dataproc / Data Fusion / BigQuery(SQL) が担当
- オーケストレーションは Cloud Composer / Workflows が担当(→ 2.3)
③ ネットワーキングの基礎
パイプラインのコンポーネント間の通信経路を設計します。
| 概念 | 説明 |
| VPC / サブネット | コンポーネントを配置する仮想ネットワーク |
| 限定公開のGoogleアクセス(Private Google Access) | 外部IPなしで内部からGoogle APIへ到達 |
| Private Service Connect / VPCピアリング | サービス/VPC間をプライベートに接続 |
| ファイアウォール ルール | 通信の許可・拒否を制御 |
- Dataflow ワーカーは外部IPを無効化し、限定公開のGoogleアクセスで通信させるのがセキュアな基本形
- オンプレ連携は Cloud VPN / Cloud Interconnect でプライベート接続
④ データ暗号化
🔑 セクション1の復習
パイプラインでも保存時・転送時とも自動で暗号化される。
| 方式 | 鍵の管理 | 使いどころ |
| デフォルト暗号化 | Google | 特別な要件がなければこれ(自動) |
| CMEK(顧客管理鍵) | あなた(Cloud KMS) | 鍵のローテーション/無効化を制御したい |
| CSEK(顧客提供鍵) | あなた(持ち込み) | 鍵をGCPに預けたくない厳格要件 |
- Pub/Sub・Dataflow・BigQuery・GCS いずれも CMEK 対応。規制要件があれば CMEK を指定
2.2 パイプラインの構築
① データクレンジング
取り込んだ生データを分析に使える品質へ整える工程。
生データ ──► [欠損処理] ──► [重複排除] ──► [型/形式統一] ──► [外れ値処理] ──► クリーンデータ
NULL補完 重複行削除 日付/数値の正規化 範囲外の除外
- 担当ツール:Dataflow(コード)/ Cloud Data Fusion(GUI、Wrangler機能)/ BigQuery(SQL) / Dataprep
- デッドレターパターン:処理できない不正レコードを別の場所(デッドレターキュー/テーブル)へ隔離し、パイプラインを止めない
② サービスの見極め(最重要)
データ処理の中核サービスを役割で覚えます。
| サービス | 種別 | ひとことで | 主な使いどころ |
| Pub/Sub | メッセージング | データを「受け取る」グローバルなキュー | ストリーミング取り込み・イベント連携・バッファ |
| Dataflow | 処理(サーバーレス) | Apache Beamで「処理する」 | バッチ/ストリーミング統一・自動スケール |
| Apache Beam | プログラミングモデル | Dataflowの「中身」 | 統一パイプライン記述(移植可能) |
| Dataproc | 処理(マネージドクラスタ) | Spark/Hadoopを「動かす」 | 既存Spark/Hadoop資産の移行 |
| Cloud Data Fusion | 処理(GUI) | ノーコードで「組み立てる」 | GUIでETL/ELT・コードを書かない |
| BigQuery | DWH/処理 | SQLで「変換・分析する」 | ELT・大規模分析・BQML |
| Apache Kafka | メッセージング(OSS) | OSS版の「受け取る」 | 既存Kafka資産(Managed Service for Kafka) |
【役割の動詞で覚える】
Pub/Sub = 受け取る(緩衝・配信)
Dataflow = 処理する(サーバーレスで柔軟)
Dataproc = (Sparkを)動かす(既存資産の移行)
Data Fusion= 組み立てる(GUIノーコード)
BigQuery = 変換・分析する(SQL)
Pub/Sub の役割
グローバルなメッセージング基盤。送信側(パブリッシャー)と受信側(サブスクライバー)を疎結合にする。
パブリッシャー ──► [トピック] ──► [サブスクリプション] ──► サブスクライバー
(アプリ/IoT) (Dataflow等)
メッセージを保持し、配信を仲介
- トピック (Topic):メッセージの送り先(カテゴリ)
- サブスクリプション (Subscription):トピックからメッセージを受け取る購読口
- プル型 (Pull):受信側が取りに行く
- プッシュ型 (Push):Pub/SubがHTTPエンドポイントへ送る
- 役割:急なデータ流入を吸収するバッファ、複数受信者へのファンアウト、送受信の疎結合化
- 配信は少なくとも1回(at-least-once)が基本(重複あり得る)。Dataflowと組むと Exactly-once を実現できる
Dataflow(Apache Beam)の役割
サーバーレスでバッチもストリーミングも同じコードで処理できる。Apache Beam という統一モデルを実行する。
【Apache Beam の基本要素】
PCollection ── データの集合(バッチ=有界 / ストリーム=無界)
PTransform ── データへの変換操作(Map, Filter, GroupBy, Window…)
Pipeline ── PTransform をつないだ処理フロー全体
Runner ── 実行エンジン(Dataflow Runner で GCP上で動く)
入力 ─► [Read] ─► [Transform] ─► [Transform] ─► [Write] ─► 出力
PCollection を PTransform で次々と変換していく
- 有界データ (Bounded) = バッチ(ファイル等、終わりがある)
- 無界データ (Unbounded) = ストリーミング(Pub/Sub等、終わりがない)
- 同じBeamコードで両方書けるのが最大の特徴("Batch and Streaming Unified")
- サーバーレス:クラスタ管理不要、ワーカーを自動スケール
Dataproc の役割
マネージドな Spark / Hadoop クラスタ。OSSのビッグデータエコシステムをそのまま動かす。
【Hadoop/Sparkエコシステム】
Spark ── 高速な分散処理エンジン(インメモリ)
Hadoop ── 分散処理の基盤(MapReduce, YARN)
HDFS ── 分散ファイルシステム(GCPではGCSに置き換え推奨)
Hive/Pig/Presto ── SQL/スクリプトでの処理
- 使いどころ:既存の Spark/Hadoop のコード資産をそのまま移行したいとき
- クラスタは数分で起動。ジョブ完了後に破棄する ephemeral(一時)クラスタ+GCS が低コストの定石
- Dataflow がサーバーレスなのに対し、Dataproc はクラスタを意識する
Cloud Data Fusion の役割
GUIでドラッグ&ドロップしてパイプラインを組むノーコード/ローコードETL(CDAPベース)。
- コードを書かずに、視覚的にソース→変換→シンクをつなぐ
- Wrangler:データを見ながら対話的にクレンジング・整形
- 内部では Dataproc 上で Spark として実行される
- 使いどころ:コードを書きたくない/書けないチーム、迅速なETL構築
③ 変換(Transformations)
バッチ処理 vs ストリーミング処理
| 観点 | バッチ処理 | ストリーミング処理 |
| データ | 蓄積済み(有界) | 到着し続ける(無界) |
| 重視 | スループット(量) | レイテンシ(速さ) |
| 例 | 日次の売上集計、月次レポート | リアルタイム不正検知、ライブ ダッシュボード |
| 代表サービス | Dataflow(バッチ) / Dataproc / BigQuery | Dataflow(ストリーミング) / Pub/Sub |
| 結果のタイミング | まとめて一度に | データ到着ごと/ウィンドウごと |
バッチ : [====蓄積====] ──► 一括処理 ──► 結果
ストリーミング: ● ● ● ● ● ● ● ● …(途切れず流れる)──► 逐次処理 ──► 逐次結果
ストリーミングのウィンドウ処理(Windowing)
無界ストリームは「いつ集計するか」を決められないため、時間で区切る=ウィンドウを使う。これは試験頻出。
① 固定ウィンドウ(Fixed / Tumbling)── 重ならない一定幅
|--w1--|--w2--|--w3--| 例: 1分ごとの件数
隙間なく、重なりなく分割
② スライディングウィンドウ(Sliding / Hopping)── 重なる
|---w1---|
|---w2---|
|---w3---| 例: 直近5分の平均を1分ごとに算出
幅(window) と 間隔(period) を指定。重複あり
③ セッションウィンドウ(Session)── 活動の途切れで区切る
●●● ___gap___ ●●●● __gap__ ●● 例: ユーザーの操作セッション
一定時間(gap)アクティビティが無ければ区切る。幅は可変
| ウィンドウ | 重なり | 幅 | 典型用途 |
| 固定 (Fixed/Tumbling) | なし | 固定 | 1分ごとの集計、定期的なバケット集計 |
| スライディング (Sliding/Hopping) | あり | 固定 | 移動平均、直近N分の傾向 |
| セッション (Session) | なし | 可変 | ユーザーセッション、操作のまとまり |
💡 補足
グローバルウィンドウ (Global) はストリーム全体を1つの窓とみなす(カスタムトリガー必須)。
ウォーターマーク(Watermark)と遅延到着データ(Late data)
イベント時刻 ──► ネットワーク遅延・順序逆転 ──► 処理時刻
12:00:30 12:00:35 に到着
- イベント時刻 (Event time):データが実際に発生した時刻
- 処理時刻 (Processing time):システムが処理した時刻
- ネットワーク遅延などで、データは順序が乱れて遅れて届くことがある
- ウォーターマーク (Watermark):「この時刻までのデータはもう出揃った」とシステムが判断する境界線。これを過ぎたウィンドウは集計を確定できる
- 遅延到着データ (Late data):ウォーターマークより後に届いたデータ
- 許容遅延 (allowed lateness) を設定すると、その範囲の遅延データを後から取り込み再計算できる
- 許容を超えた超遅延データは破棄(またはデッドレターへ)
トリガー(Trigger)
- トリガー = 「ウィンドウの集計結果をいつ出力するか」を制御する仕組み
- 種類:
- イベント時刻トリガー:ウォーターマーク到達時に出力(既定)
- 処理時刻トリガー:一定の処理時間ごとに出力
- データ駆動トリガー:N件たまったら出力
- 複合トリガー:早期結果(途中経過)+遅延データでの再出力を組み合わせる
- 早期発火+遅延更新:ウォーターマーク前に途中経過を出し、遅延データが来たら結果を更新、という運用ができる
Exactly-once 処理
- Exactly-once:各レコードを重複なく・漏れなく ちょうど1回処理する保証
- Dataflow は内部でチェックポイントと重複排除を行い、Pub/Sub と組み合わせて Exactly-once を実現できる
- 比較:at-least-once(重複あり得る)/ at-most-once(欠損あり得る)/ exactly-once(ちょうど1回)
④ 処理ロジック
- フィルタリング(条件で絞る)、マッピング(1→1変換)、集約(GroupBy/Combine)、結合(Join)、エンリッチメント(外部データで補強)
- Beam では
ParDo(要素ごとの処理)、GroupByKey、Combine、Flatten、Partition などの PTransform を組み合わせる
- サイド入力(Side Input):メインのストリームに参照データ(マスタ等)を合流させる仕組み
⑤ AI によるデータエンリッチメント
近年追加された範囲。パイプライン内で AI/ML を使ってデータを補強する。
| 手法 | 説明 | 例 |
| BigQuery ML の生成AI関数 | SQLからLLM/モデルを呼び出す | ML.GENERATE_TEXT(要約・分類)、ML.GENERATE_EMBEDDING(ベクトル化) |
| Vertex AI 連携 | パイプラインからモデル推論を呼ぶ | 画像分類、感情分析、エンティティ抽出 |
| Cloud DLP | PII検出・マスキング | 取り込み時に機密情報を秘匿 |
- 例:レビューテキストを
ML.GENERATE_TEXT で感情分類して列を追加、画像をVertex AIで分類してタグ付け
⑥ データの取得とインポート
| 方法 | 用途 |
| bq load / BigQuery への直接ロード | GCS/ローカルのファイルをBigQueryへ取り込み |
| Pub/Sub → Dataflow → BigQuery | ストリーミング取り込みの定番 |
| Storage Transfer Service | 他クラウド/オンプレ → GCS の大容量転送 |
| BigQuery Data Transfer Service | SaaS/他DWHからの定期取り込み |
| Datastream | DBの変更をCDCで継続複製 |
| Dataflow テンプレート | Googleが用意した定型パイプライン(GCS→BQ等)をパラメータ実行 |
- 形式:CSV / JSON / Avro / Parquet / ORC。Avro/Parquet は型情報を持ち高速・堅牢
- BigQuery 外部テーブル / BigLake:ロードせずGCS上のファイルを直接クエリ
⑦ 新規データソースとの統合
- 新しいSaaSやDBを追加する際、コネクタ(Data Fusion のプラグイン、Datastream の対応ソース等)を確認
- スキーマの違いを吸収する変換層を設け、既存パイプラインに合流させる
2.3 パイプラインのデプロイと運用化
① ジョブの自動化とオーケストレーション
複数のジョブを「順番に・依存関係を守って・定期的に」実行する。代表は2つ。
Cloud Composer(マネージド Apache Airflow)
【DAG = 有向非巡回グラフ】タスクの依存関係をコードで定義
┌──────────┐
│ extract │
└────┬─────┘
▼
┌──────────┐
│transform │
└────┬─────┘
┌────┴─────┐
▼ ▼
┌────────┐ ┌────────┐
│load_bq │ │load_gcs│ ← 並列に実行
└────┬───┘ └───┬────┘
└────┬────┘
▼
┌──────────┐
│ notify │
└──────────┘
- DAG(有向非巡回グラフ):タスクと依存関係をPythonコードで定義
- 使いどころ:複雑な依存関係、多数のタスク、リトライ・分岐・スケジュールが必要な本格的なワークフロー
- 豊富な Operator(BigQueryOperator, DataflowOperator 等)でGCPサービスを呼べる
- Airflow資産(既存DAG)をそのまま使える
Workflows
- サーバーレスなAPIオーケストレーション。YAML/JSONでステップを定義
- 使いどころ:軽量なサービス連携、HTTP/API・Cloud Functions・Cloud Runの順次/分岐実行
- イベント駆動・低コスト・インフラ管理不要。シンプルな連携向き
【使い分けの第一感】
複雑な依存・多数タスク・データパイプライン → Cloud Composer (Airflow)
軽量なAPI/サービス連携・イベント駆動・低コスト → Workflows
- 補助:Cloud Scheduler(cronで定期トリガー)、Eventarc(イベント駆動の起動)
② CI/CD(継続的インテグレーション・デプロイ)
パイプラインのコードを自動でテスト・ビルド・デプロイする。
コード push ──► [Cloud Build トリガー] ──► テスト ──► ビルド ──► デプロイ
(Git) (ソースの変更を検知) 単体/結合 成果物 Dataflowテンプレート/
Composer DAG 等を配置
- Cloud Build:ビルド・テスト・デプロイを自動化するCI/CDサービス(
cloudbuild.yaml で定義)
- 使いどころ:Dataflow テンプレートのビルド&登録、Composer の DAG を GCS の dags フォルダへ配置、Terraformでのインフラ適用
- Artifact Registry:ビルドしたコンテナ/パッケージの保管庫
- ベストプラクティス:dev→staging→prod の環境ごとに自動デプロイし、テストを通過したものだけ昇格
📌 このセクションの基礎まとめ
- パイプラインは ソース→取り込み→変換→シンク + オーケストレーション+CI/CD
- 処理サービスは役割で覚える:Pub/Sub=受け取る、Dataflow=処理する(サーバーレス)、Dataproc=Sparkを動かす、Data Fusion=GUIで組む、BigQuery=SQLで変換
- ストリーミングの核:ウィンドウ(固定/スライディング/セッション)・ウォーターマーク・トリガー・遅延データ・Exactly-once
- オーケストレーション:複雑な依存=Composer、軽量連携=Workflows
- 自動デプロイは Cloud Build による CI/CD
🔧 実践(中堅) 🎯 発展(シニア)
💡 このレベルのゴール
「どう設計判断するか」「トレードオフは何か」「試験のひっかけ」に焦点を当てます。基礎概念は「基礎」タブを参照。
最重要セクションのため、判断軸とひっかけを特に厚く扱います。
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
- BigQuery への書き込みは Storage Write API で Exactly-once セマンティクスを利用可能
⚠️ ひっかけ
「重複なくちょうど1回処理したい」→ Dataflow の Exactly-once 処理(Pub/Subの重複をDataflowが吸収)
🔧🎯 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
🎯 要点と暗記
💡 使い方
試験直前の総ざらい用。暗記すべき表と一問一答でセクション2を固めます。
最重要セクション(~25%)。ここの取りこぼしは致命的なので確実に。
🔑 暗記必須テーブル
処理サービスの役割(動詞で覚える)
| サービス | 役割(動詞) | 選ぶ条件 |
| Pub/Sub | 受け取る(緩衝・配信) | ストリーミング取り込み・疎結合・バッファ |
| Dataflow | 処理する(サーバーレス) | Beam・ストリーミング/バッチ統一・自動スケール |
| Dataproc | (Sparkを)動かす | 既存Spark/Hadoop資産の移行 |
| Cloud Data Fusion | GUIで組み立てる | ノーコード/ローコードETL |
| BigQuery | SQLで変換・分析 | ELT・大規模分析・BQML |
Dataflow vs Dataproc vs Data Fusion
| キーワード | 答え |
| サーバーレス / Apache Beam / ストリーミング | Dataflow |
| 既存Spark/Hadoop/Hive / リフト&シフト | Dataproc |
| GUI / ノーコード / ドラッグ&ドロップ | Cloud Data Fusion |
| SQLだけ / BigQuery内ELT | BigQuery / Dataform |
ストリーミングのウィンドウ
| ウィンドウ | 重なり | 幅 | 用途 |
| 固定 (Fixed/Tumbling) | なし | 固定 | 毎分/毎時の定期集計 |
| スライディング (Sliding/Hopping) | あり | 固定 | 移動平均・直近N分 |
| セッション (Session) | なし | 可変 | 無活動gapで区切る操作のまとまり |
ストリーミング処理の核概念
| 用語 | 一言 |
| イベント時刻 | データが発生した時刻(集計はこれが原則) |
| 処理時刻 | システムが処理した時刻 |
| ウォーターマーク | 「この時刻まで出揃った」と見なす境界 |
| 遅延到着データ | ウォーターマーク後に届くデータ(許容遅延で再計算) |
| トリガー | ウィンドウ結果を「いつ出すか」の制御 |
| Exactly-once | 重複なく漏れなくちょうど1回処理 |
Dataflow のスケール/安定化機能
| 機能 | 対象 | 効果 |
| Dataflow Shuffle | バッチ | シャッフルをサービス側へ。VM負荷軽減 |
| Streaming Engine | ストリーミング | 状態をサービス側へ。オートスケール高速・安定 |
| FlexRS | バッチ | 遅延スケジューリングで低コスト |
Pub/Sub の重要機能
| 機能 | 一言 |
| 順序指定キー | 同じキー内で順序保証(全順序ではない) |
| デッドレタートピック | 失敗メッセージを隔離・パイプライン停止防止 |
| シーク | 過去時刻/スナップショットへ再生位置を戻す |
| メッセージ保持 | 確認後も保持しシークでの再処理を可能に |
| Pub/Sub Lite | ゾーン単位・低コスト・要キャパシティ計画 |
オーケストレーション
| キーワード | 答え |
| DAG / Airflow / 複雑な依存 / 多数のデータジョブ | Cloud Composer |
| 軽量 / サーバーレス / API・HTTP連携 / 低コスト | Workflows |
| 単純な定期トリガー(cron)だけ | Cloud Scheduler |
| CI/CD・ビルド・テスト・デプロイ自動化 | Cloud Build |
✅ 一問一答(確認テスト)
クリックで解答を表示。全15問を即答できるか確認しましょう。
Q1. 既存のApache Spark/HadoopジョブをほぼそのままGCPへ移行したい。最適なサービスは?
A. Dataproc(マネージドSpark/Hadoop。Dataflowに書き換える必要がない。コスト最適化はephemeralクラスタ+GCS)
Q2. サーバーレスで、バッチとストリーミングを同じコードで処理したい。使うサービスとモデルは?
A. Dataflow(実行)+ Apache Beam(プログラミングモデル)
Q3. 「直近5分間の平均を1分ごとに更新」する集計に使うウィンドウは?
A. スライディングウィンドウ(Sliding/Hopping)(window=5分, period=1分)
Q4. 「ユーザーが30分操作しなければセッション終了」とみなす集計のウィンドウは?
A. セッションウィンドウ(Session)(gap=30分。幅は可変)
Q5. ネットワーク遅延でイベントが順不同・遅れて届くが、正しく時間集計したい。基準とする時刻と仕組みは?
A. イベント時刻(Event time)+ ウォーターマーク(Watermark)。遅延データは許容遅延(allowed lateness)とトリガーで反映する。
Q6. Pub/Subで処理に失敗し続けるメッセージがパイプラインを詰まらせている。対策は?
A. デッドレタートピック(DLQ)(規定回数失敗したメッセージを退避し、本処理を継続)
Q7. バグ修正後、過去に配信済みのメッセージを再処理したい。Pub/Subの機能は?
A. シーク(Seek)(過去のタイムスタンプ/スナップショットへ再生位置を戻す。メッセージ保持が前提)
Q8. メッセージを重複なく・漏れなくちょうど1回処理したい。どう実現する?
A. Dataflow の Exactly-once 処理(Pub/Subは at-least-once で重複し得るが、Dataflowが重複排除・チェックポイントで吸収。BigQuery書き込みは Storage Write API)
Q9. ストリーミングのオートスケールを高速・滑らかにし、ワーカーの安定性を高めたい。有効化する機能は?
A. Streaming Engine(ウィンドウ状態をサービス側へ移し、ワーカーを軽量化)。バッチのシャッフル負荷なら Dataflow Shuffle。
Q10. 多数のジョブを複雑な依存関係に従って実行・スケジュールしたい。使うサービスは?
A. Cloud Composer(マネージドAirflow、DAG)。軽量なAPI連携なら Workflows、単純な定期実行だけなら Cloud Scheduler。
Q11. コストを最優先し、スループットが予測でき、ゾーン単位の可用性で良いメッセージング基盤は?
A. Pub/Sub Lite(事前にキャパシティを予約。標準のPub/Subより大幅に安価)
Q12. パイプラインのコードをGit管理し、テスト・ビルド・デプロイを自動化して再現可能にしたい。中心となるサービスは?
A. Cloud Build(CI/CD)。Dataflowはテンプレート化しArtifact Registryへ、ComposerのDAGはGCSのdags/へ自動デプロイ。dev→staging→prodで昇格。
Q13. アナリストがコードを書かずにGUIでETLパイプラインを組みたい。最適なサービスは?
A. Cloud Data Fusion(ドラッグ&ドロップ。Wranglerで対話的クレンジング。内部はDataproc/Spark)
Q14. 同一ユーザーのイベントだけは到着順序を保ったまま処理したい。Pub/Subで使う機能は?
A. 順序指定キー(Ordering Key)(キー単位で順序保証。全メッセージのグローバル順序ではない点に注意)
Q15. BigQuery内のテキスト列をSQLだけで要約・分類してエンリッチしたい。使うものは?
A. BigQuery ML の生成AI関数 ML.GENERATE_TEXT(埋め込み生成なら ML.GENERATE_EMBEDDING)
🎯 ひっかけ注意ポイント
❌ ひっかけ注意
- 「既存Spark/Hadoop」と出たら Dataproc(Dataflowに誘導する選択肢は罠)
- 「サーバーレス・Beam・ストリーミング統一」は Dataflow
- 「GUI/ノーコード」は Cloud Data Fusion 一択
- Composer と Workflows を混同しない:複雑な依存=Composer、軽量API連携=Workflows
- ウィンドウの取り違え:移動平均=スライディング、無活動で区切る=セッション、定期=固定
- 処理時刻で集計する選択肢は罠 → 正しく集計するならイベント時刻+ウォーターマーク
- Pub/Sub は既定で at-least-once(重複あり得る)。「重複なく1回」は Dataflowの Exactly-once
- Shuffle(バッチ)と Streaming Engine(ストリーミング) を取り違えない
- Pub/Sub Lite は「コスト最優先+キャパシティ計画可」のときだけ。可用性・無設定重視なら標準Pub/Sub
- 「デッドレター=隔離」「シーク=巻き戻し再処理」を逆に覚えない
- Cloud Scheduler は起動トリガーであり、依存管理はできない(依存はComposer)
📝 セルフチェック
📍 ナビゲーション