03_要点と暗記

セクション2 要点と暗記:データの取り込みと処理 🎯

試験直前の総ざらい用。暗記すべき表一問一答でセクション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

✅ 一問一答(確認テスト)

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


🎯 ひっかけ注意ポイント


📝 セルフチェック