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 のアーキテクチャ: クラスタ、サーバーレス、ストレージ、メタストア

    --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}
```
      --conf spark.sql.autoBroadcastJoinThreshold=64m
      ```

    - ホットなパーティションにはキーにソルトを追加します。マップ側で事前集計を適用し、早期にフィルタリングします。
    - Adaptive Query Execution (AQE)を有効にして、シャッフル後のパーティションを結合し、スキューのあるJoinを処理します:
  --conf spark.sql.adaptive.enabled=true
  ```
      --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 (...)
```

セキュリティ、オブザーバビリティ、コスト

    --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
 ```
  1. セキュリティとネットワーキングを強化する

    • GCSパス、メタストア、BigQueryデータセットに必要なロールのみを付与した専用のサービスアカウントでクラスタを実行します。制限されたサブネットにPrivate Google Accessを持つプライベートIPクラスタを作成し、ファイアウォールルールを介してUIアクセスを制限します。
    • 理由:最小権限とネットワーク分離は攻撃対象領域を縮小します。プライベートなコントロールプレーンの下り(egress)は、パブリックへの露出を回避します。
  2. ロギング、履歴、アラートを整備する

    • GCSへのSparkイベントログを有効にし、History Serverをデプロイします。ドライバ/YARNログを保持期間付きでCloud Loggingにルーティングします。長時間の保留コンテナ、繰り返されるタスクの失敗、または過剰なジョブ期間に対してMonitoringアラートを追加します。
    • 理由:一元化されたログは根本原因分析をサポートします。プロアクティブなアラートは、スキュー、OOM、またはI/Oの劣化を早期に検出します。
  3. アドホックおよび弾力的なスパイク対応のために、Dataproc Serverlessで選択的にモダナイズする

    • 散発的または探索的なSpark SQLワークロードをDataproc Serverlessに移行します。夜間のパイプラインは、サーバーレスで完全に検証されるまでエフェメラルクラスタで維持します。
    • 理由:サーバーレスはクラスタ運用を不要にし、自動的にスケールするため、予測不能な負荷に最適です。既存のワークフローは最小限のコード変更で継続できます。
  4. オブジェクトストアのコミッターと小ファイル管理を検証する

    • 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.

試験に合格する →

Googleを閲覧 →

Related guides

オールインワンアクセス

1つのサブスクリプション。すべての試験。

すべてのプランで、無制限の回答検索、模擬試験、AI解説、および完全なリソースライブラリが利用可能 — 20以上の言語に対応。

月額
24.87
Just €0.83/day
すべて含まれています:
  • 無制限の回答検索
  • 無制限の模擬試験
  • AIを活用した解説
  • 完全なリソースライブラリ
  • 20以上の言語
  • 毎週のコンテンツ更新
  • 特典 & 紹介
  • 優先サポート
無料トライアルを開始

クレジットカード不要*

ベストバリュー
12ヶ月
179.87
Just €0.49/daySave 40%
すべて含まれています:
  • 無制限の回答検索
  • 無制限の模擬試験
  • AIを活用した解説
  • 完全なリソースライブラリ
  • 20以上の言語
  • 毎週のコンテンツ更新
  • 特典 & 紹介
  • 優先サポート
無料トライアルを開始

クレジットカード不要*

✓ 無料プランが含まれています · ✓ いつでもキャンセル可能 · ✓ すべてのプランで製品の全機能が利用可能