Google PDE: Spark、Dataproc、分散データ処理 — 学習ガイド
こちらの一部です: Google Professional Data Engineer — 学習ガイド. 検証済みの解答で練習: Google試験ハブ, または時間制限付き模擬試験に挑戦: ExamRoll.io.
概要
Google Cloud Dataproc 上の Apache Spark は、分散データ処理のためのマネージドで伸縮自在なプラットフォームを提供します。制御の要件、ランタイムの変動性、管理オーバーヘッドに応じて、実行時間の長い (long-running) またはエフェメラルな Dataproc クラスタ、あるいは Dataproc Serverless for Spark を選択できます。Spark は、回復力のある抽象化 (RDD)、リレーショナル API (DataFrame と Spark SQL)、そして大規模な反復処理やバッチ ETL に最適化されたフォールトトレラントな DAG 実行エンジンを提供します。Google Cloud では、Cloud Storage が HDFS に代わって永続的で低コストなストレージを提供し、BigQuery コネクタが直接的な分析オフロードを可能にし、Dataproc Metastore がスキーマ管理を一元化します。効果的なソリューションは、ストレージとコンピューティングのライフサイクルを整合させ、ワークロードに合わせて Spark をチューニングし、オブザーバビリティを実装し、最小権限とネットワーク分離によるセキュリティを適用します。
Dataproc のアーキテクチャ: クラスタ、サーバーレス、ストレージ、メタストア
- クラスタタイプとノードの役割
- プライマリ (マスター) ノードは、YARN、HDFS NameNode (使用する場合)、Spark ドライバの UI をホストします。HA モードでは複数のプライマリノードを使用します。
- ワーカーノードは、エグゼキュータと HDFS DataNode (使用する場合) を実行します。
- セカンダリ/補助ワーカーは、通常、プリエンプティブル/スポット VM であり、HDFS の役割を持たずに、伸縮性のある低コストのキャパシティを提供します。
- イメージは OS とコンポーネントのバージョン (例: 2.1-debian11, 2.2-ubuntu20) をバンドルします。イメージバージョンを固定することで、Spark/Hadoop の互換性を制御し、意図的にアップグレードを行います。
- Component Gateway は、UI (Spark History Server, YARN RM) を HTTPS 経由で安全に公開します。
- Dataproc Serverless for Spark
- クラスタのプロビジョニングが不要で、自動オートスケーリングが行われ、エグゼキュータとドライバに対して秒単位の課金となります。散発的またはバースト的なジョブ、あるいは運用オーバーヘッドを最小限に抑えたい場合に最適です。
- トレードオフ: クラスタよりも低レベルの調整項目が少ない、ジョブの開始レイテンシがウォーム状態のクラスタより高くなる可能性がある、トラブルシューティングにはサーバーレスのメトリクスとイベントログを使用する、などがあります。
- オートスケーリング
- クラスタのオートスケーリングポリシーは、YARN/Spark のメトリクスとクールダウン期間に基づいてワーカーを追加/削除し、プライマリとセカンダリのワーカーグループを個別に調整します。
- サーバーレスのオートスケーリングはサービスによって管理されます。最適なスケーリングのためには、パーティション並列になるように設計し、直列処理のボトルネックを避ける必要があります。
- ストレージとコネクタ
- system-of-record (信頼できる唯一の情報源) としては Google Cloud Storage (GCS) を推奨します。GCS はコンピューティングとストレージを分離し、永続ディスクのコストを削減し、クラスタのライフサイクルを超えて存続します。
- GCS コネクタ (gs://) は Hadoop/Spark と統合されます。オブジェクトストレージへの書き込みにはコミットプロトコルが使用されます。GCS でのリネームのオーバーヘッドを削減し、ジョブのコミットを高速化するために、FileOutputCommitter アルゴリズム v2 を設定します:
--conf mapreduce.fileoutputcommitter.algorithm.version=2
```
- カラムプルーニング (列の枝刈り) とプレディケートプッシュダウン (述語プッシュダウン) を備えた Parquet/ORC を使用します。小さなファイルはコンパクション (統合) によって管理し、効率的なスキャンのためにファイルあたり 128~512 MiB を目標にします。
- Hive メタストア
- スキーマとテーブルのメタデータは、Dataproc Metastore (マネージド Apache Hive Metastore) または Cloud SQL をバックエンドとするメタストアに一元化し、クラスタ間でカタログを共有します。
- 永続性のために GCS を指す外部テーブルを使用します。スキャンコストを抑制するために、日付/時間でパーティション分割します。
- ジョブ、初期化、ワークフロー
- spark, pyspark, spark-sql, または hadoop ジョブをサブミットします。初期化アクションは、クラスタ作成時に追加のライブラリやエージェント (例: コネクタ、Python ライブラリ) をインストールします。
- ワークフローテンプレートは、複数ステップのパイプラインをパラメータ化します。ワークフローごとにエフェメラルクラスタを作成し、終了後に破棄することができます。これにより、分離性が向上し、アイドルコストが削減されます。
- バッチ ETL にはエフェメラルクラスタが推奨されます。データとメタストアはクラスタの外部 (GCS, Dataproc Metastore, BigQuery) に配置します。
- BigQuery との統合
- Spark BigQuery コネクタは BigQuery を直接読み書きします。スループットのためには BigQuery Storage Read API を、低レイテンシで exactly-once のストリーミング挿入のためには Write API の使用を検討します。
- テーブルのメンテナンスでは、BigQuery で下流の MERGE/パーティション上書きを実行し、ロードをアトミックに完了させます。
### Sparkモデル、パフォーマンスチューニング、信頼性
- APIと実行
- RDD: 低レベル、イミュータブル、Scala/Javaでは型安全。パーティショニングと永続化を制御します。
- DataFrame/Dataset: リレーショナル、Catalystで最適化済み。クエリ最適化とコード生成のため、ETLにはこちらを推奨します。
- トランスフォーメーションは遅延評価(map, filter, join)です。アクションが実行をトリガーします(count, collect, save)。Sparkはシャッフルによって分割されたステージのDAGを構築し、タスクはパーティションごとに実行されます。
- パーティショニングとシャッフル
- 入力パーティショニング: 全てのコアを活用するのに十分なパーティション数が必要です。エグゼキュータの総コア数の2〜4倍から始めます。`spark.default.parallelism`(RDD用)やリーダーオプション(DataFrame用)で制御します。
- シャッフルパーティション: デフォルトの200では、過小または過剰なプロビジョニングになりがちです。調整方法:
--conf spark.sql.shuffle.partitions= {total_executor_cores * 2 to 3}
```
- ワイドトランスフォーメーション後のパーティションあたり約100〜256 MiBを目標とします。小さすぎるとスケジューラのオーバーヘッドが発生し、大きすぎるとエグゼキュータのOOMリスクが高まります。
- シャッフルはjoins, groupBy, orderByで最もコストがかかる部分です。エグゼキュータのメモリとディスクを十分に確保し、クラスターでシャッフルが多発する場合はローカルSSDを検討します。
- スキュー(偏り)とJoin戦略
- スキューを検出します(ロングテールのタスク実行時間、大きなパーティションサイズ)。緩和策:
- シャッフルを避けるために小さなテーブルをブロードキャストします:
- スキューを検出します(ロングテールのタスク実行時間、大きなパーティションサイズ)。緩和策:
--conf spark.sql.autoBroadcastJoinThreshold=64m
```
- ホットなパーティションにはキーにソルトを追加します。マップ側で事前集計を適用し、早期にフィルタリングします。
- Adaptive Query Execution (AQE)を有効にして、シャッフル後のパーティションを結合し、スキューのあるJoinを処理します:
--conf spark.sql.adaptive.enabled=true
```
- キャッシング、チェックポインティング、リネージ
- 再利用されるホットな中間DataFrameは慎重にキャッシュします。OOMを避けるため、
MEMORY_AND_DISKを推奨します。 - 長いリネージを持つ場合はGCSやHDFSにチェックポイントを作成し、障害発生時の再計算範囲を限定します。
- 再利用されるホットな中間DataFrameは慎重にキャッシュします。OOMを避けるため、
- エグゼキュータと動的割り当て
- 並列処理とGCオーバーヘッドのバランスを取るために、エグゼキュータのサイズを適正化します:
- エグゼキュータあたりのコア数: I/OとCPUのバランスが取れたタスクには2〜5コア。コア数を減らすとGCの一時停止が短縮されます。
- メモリオーバーヘッド: ワイドシャッフルには
spark.yarn.executor.memoryOverheadを設定します。 - クラスターでは、External Shuffle Serviceと共に動的割り当てを有効にし、ワークロードに応じてエグゼキュータをスケールさせます:
- 並列処理とGCオーバーヘッドのバランスを取るために、エグゼキュータのサイズを適正化します:
--conf spark.dynamicAllocation.enabled=true
--conf spark.shuffle.service.enabled=true
--conf spark.dynamicAllocation.minExecutors=0
--conf spark.dynamicAllocation.maxExecutors=200
```
- バッチETLのための耐障害性パターン
- べき等な書き込み: 一時/ステージングパスに書き込み、ディレクトリレベルのコミットでアトミックに昇格させます。BigQueryの場合は、ステージングテーブルに書き込み、`MERGE`を使用します:
MERGE target t USING staging s
ON t.id = s.id
WHEN MATCHED THEN UPDATE SET ...
WHEN NOT MATCHED THEN INSERT (...)
```
- 増分処理:
ingestion_dateパーティションに対してウォーターマークベースのフィルタリングを使用します。再処理を避けるため、処理済みマニフェストをGCSで管理します。 - デッドレター処理: パース/バリデーションエラーが発生した場合、不正なレコードを診断情報と共に隔離パス/テーブルに分岐させます。厳密なスキーマ適用と組み込みのDLQが必要な場合はDataflowを検討します。Sparkでは、レコードごとのtry/catchと別のシンクを実装します。
セキュリティ、オブザーバビリティ、コスト
- IDとアクセス
- 専用のサービスアカウントと最小権限のIAMを使用してクラスタとジョブを実行します。必要なロールのみを付与します。例:
- インスタンスのサービスアカウントに roles/dataproc.worker
- GCSのI/Oパスに roles/storage.objectViewer または objectAdmin
- ターゲットデータセットに roles/bigquery.dataEditor
- Dataproc Serverlessでは、ジョブごとのサービスアカウントを使用してアクセス範囲を限定します。
- 専用のサービスアカウントと最小権限のIAMを使用してクラスタとジョブを実行します。必要なロールのみを付与します。例:
- ネットワークの分離と暗号化
- VPCサブネット内でプライベートIPクラスタを使用し、マスターのUIをファイアウォールで制限し、パブリックな下り(egress)なしでGCS/BigQueryにアクセスするためにPrivate Google Accessを有効にします。
- 一元管理のためにShared VPCプロジェクトにクラスタを配置します。オプションで、クラスタ内認証のためにDataprocでKerberosを有効にします。
- CMEKで保存データを暗号化します。GCSバケット、Persistent Disks、Dataproc Metastore、BigQueryでCMEKを構成します。転送中のデータはデフォルトでTLSを使用します。
- ロギング、履歴、メトリクス
- GCSへのSparkイベントログを有効にし、History Serverをデプロイします:
--conf spark.eventLog.enabled=true
--conf spark.eventLog.dir=gs://bucket/spark-events/
```
- DataprocはドライバとYARNのログをCloud Loggingにストリーミングします。保持/フォレンジックのためにシンクにエクスポートします。
- Cloud Monitoringのメトリクスで監視します:YARNの保留中コンテナ、CPU、メモリ、HDFSのヘルス(使用している場合)、GCSスループット。長引くステージのリトライ、エグゼキュータの損失、投機的実行の急増に対してアラートを設定します。
- 障害分析:一般的な原因には、スキューによるストラングラー(straggler)、シャッフル中のエグゼキュータのOOM、オブジェクトストアのコミット失敗、プリエンプティブル/スポットノードの損失などがあります。リトライ回数は慎重に増やしてください。過剰なリトライはコストと遅延を増大させる可能性があります。
- コスト最適化
- エフェメラルクラスタまたはDataproc Serverlessを使用してアイドルコストを回避します。永続ディスクを最小限に抑えるためにデータはGCSに保持します。
- ピーク需要を吸収するためにプリエンプティブル/スポットのセカンダリワーカーを追加します。失われたノードのタスクはリトライされるため、再計算を前提に設計します。マスターはプリエンプティブルノードに配置しないでください。
- マシンタイプを適正化し、キューが空のときにはオートスケーリングを使用してキャパシティを縮小します。スキャンコストとCPUを削減するために、パーティションプルーニングを伴うParquet/ORCを優先します。
- 出力をコンパクションして小さなファイルを避けます。ファイル数を減らしてサイズを大きくすると、メタデータのオーバーヘッドとジョブの実行時間が削減されます。
- 短時間の定期的なジョブ(例:週に1回の30分間のSpark ETL)では、プリエンプティブルワーカーまたはサーバーレスが最良のコストプロファイルを提供することが多いです。
#### 実践的な問題シナリオ
Acme Retail社は、夜間にSparkとHiveのETLを実行し、下流の分析システムにデータを供給している30ノードのオンプレミスHadoopクラスタを移行しようとしています。彼らは、既存のジョブを最小限の変更で再利用し、クラスタの常時管理を避け、クラスタのライフサイクルを超えてデータを永続化し、ストレージコストを削減したいと考えています。
アプローチ:
1) マネージドサービスにデータとメタデータを配置する
- すべての生データとキュレーション済みデータを、パーティショニング(例:dt=YYYY-MM-DD)を施したParquet形式でCloud Storageに保存します。
- 理由:GCSは耐久性が高く低コストであり、コンピューティングとストレージを分離するため、エフェメラルクラスタやサーバーレスジョブを永続ディスクなしで実行できます。パーティション化されたParquetは、述語プッシュダウンと効率的なスキャンを可能にします。
2) Dataproc Metastoreでカタログを一元化する
- HiveメタストアをDataproc Metastoreに移行します。GCSパスを参照する外部Hiveテーブルを作成し、既存のスキーマ/パーティションロジックを維持します。
- 理由:マネージドメタストアを使用すると、HA構成のMySQL/PostgreSQLインスタンスを運用することなく、複数のエフェメラルクラスタやサーバーレスジョブがテーブル定義を共有できます。
3) バッチETLにはエフェメラルDataprocクラスタを、オーケストレーションにはワークフローテンプレートを使用する
- 必要なイメージ(例:2.1-debian11)でクラスタを作成し、Sparkジョブ(spark-sqlおよびpyspark)を実行し、完了時にクラスタを削除するワークフローテンプレートを定義します。カスタムライブラリをインストールするために初期化アクションを追加します。
- 理由:エフェメラルクラスタはアイドルコストを排除し、ジョブの依存関係を分離します。ワークフローテンプレートは、再現性とパラメータ化(日付、入力パス)を提供します。
4) オートスケーリングとプリエンプティブルワーカーを有効にする
- 少数のコアワーカーグループと、より大きなプリエンプティブルセカンダリワーカーのプールを持つオートスケーリングポリシーをアタッチします。実行後すぐにスケールダウンするようにクールダウン期間を調整します。
- 理由:コアワーカーはクラスタの安定性を維持し、プリエンプティブルワーカーはシャッフルやワイド変換を低コストで吸収します。プリエンプションで失われたタスクはSpark/YARNのリトライによって処理されます。
5) Spark BigQueryコネクタを介してBigQueryと統合する
- ディメンション/ファクトのロードでは、Sparkの結果をステージング用のBigQueryテーブルに書き込み、その後MERGEステートメントを実行してターゲットをアトミックに更新します。直接の上書きが安全な場合は、パーティション上書きモードを使用してパーティションテーブルに書き込みます。
- 理由:BigQueryは大規模な分析とBIを提供します。ステージング+MERGEにより、バッチSparkからトランザクションのようなアップサートが実現され、下流の不整合が減少します。
6) パフォーマンスと信頼性のためにSparkをチューニングする
- エグゼキュータのコア数に応じてシャッフルパーティションを設定し、AQEを有効にします:
--conf spark.sql.shuffle.partitions=600
--conf spark.sql.adaptive.enabled=true
```
- 小さなディメンションにはブロードキャスト結合を使用し、安定性のために長いリネージをGCSにチェックポイントします。
- 理由:適切なパーティショニングはスキューとスケジューラのオーバーヘッドを削減します。AQEは実行時にデータプロファイルに適応します。チェックポイントは障害後の再計算の範囲を限定します。
セキュリティとネットワーキングを強化する
- GCSパス、メタストア、BigQueryデータセットに必要なロールのみを付与した専用のサービスアカウントでクラスタを実行します。制限されたサブネットにPrivate Google Accessを持つプライベートIPクラスタを作成し、ファイアウォールルールを介してUIアクセスを制限します。
- 理由:最小権限とネットワーク分離は攻撃対象領域を縮小します。プライベートなコントロールプレーンの下り(egress)は、パブリックへの露出を回避します。
ロギング、履歴、アラートを整備する
- GCSへのSparkイベントログを有効にし、History Serverをデプロイします。ドライバ/YARNログを保持期間付きでCloud Loggingにルーティングします。長時間の保留コンテナ、繰り返されるタスクの失敗、または過剰なジョブ期間に対してMonitoringアラートを追加します。
- 理由:一元化されたログは根本原因分析をサポートします。プロアクティブなアラートは、スキュー、OOM、またはI/Oの劣化を早期に検出します。
アドホックおよび弾力的なスパイク対応のために、Dataproc Serverlessで選択的にモダナイズする
- 散発的または探索的なSpark SQLワークロードをDataproc Serverlessに移行します。夜間のパイプラインは、サーバーレスで完全に検証されるまでエフェメラルクラスタで維持します。
- 理由:サーバーレスはクラスタ運用を不要にし、自動的にスケールするため、予測不能な負荷に最適です。既存のワークフローは最小限のコード変更で継続できます。
オブジェクトストアのコミッターと小ファイル管理を検証する
- FileOutputCommitterアルゴリズムv2を設定し、書き込み前にrepartition/coalesceを介して出力をファイルあたり256〜512 MiBにコンパクションします。
- 理由:オブジェクトストアにはアトミックなリネームがありません。最適化されたコミッターはコピー/リネームのオーバーヘッドを削減します。コンパクションは、パフォーマンスとコストの観点から小ファイル問題を緩和します。
この設計は、既存のSparkおよびHiveジョブを最小限のリファクタリングで再利用し、GCSでのデータの耐久性を確保し、スキーマを一元化し、セキュリティの爆発半径を抑制し、堅牢なオブザーバビリティを提供し、エフェメラルクラスタ、オートスケーリング、プリエンプティブルキャパシティ、およびサーバーレス実行の的を絞った使用を通じてコストを最適化します。
← メッセージング、イベントの取り込み、リアルタイムサービス · すべてのドメイン · データの取り込み、統合、移行 →
これらの問題を練習する → · ExamRoll.ioで時間制限付き練習 →
Pass the whole exam — not just this question
You found this answer. Get every verified question and explanation in one place, and save hours of prep. Free to start.
試験に合格する →