Google PDE: Spark、Dataproc 與分散式資料處理 — 學習指南

屬於 Google Professional Data Engineer — 學習指南. 使用經過驗證的解答練習: Google 考試中心, 或參加限時模擬考試: ExamRoll.io.

總覽

Google Cloud Dataproc 上的 Apache Spark 為分散式資料處理提供了一個託管式、具彈性的平台。您可以根據控制需求、執行時間的變動性以及管理負擔,在長時間執行的叢集、臨時性 Dataproc 叢集與 Dataproc Serverless for Spark 之間進行選擇。Spark 提供了具備彈性的抽象層 (RDDs)、關聯式 API (DataFrames 與 Spark SQL),以及一個容錯的 DAG 執行引擎,此引擎針對大規模的迭代式與批次 ETL 進行了最佳化。在 Google Cloud 上,Cloud Storage 取代 HDFS 作為耐用、低成本的儲存空間;BigQuery 連接器可直接進行分析卸載;而 Dataproc Metastore 則集中管理結構描述。有效的解決方案會對齊儲存與運算的生命週期、根據工作負載調整 Spark、建構可觀測性,並透過最小權限和網路隔離來實施安全性。

Dataproc 架構:叢集、Serverless、儲存與 Metastore

    --conf mapreduce.fileoutputcommitter.algorithm.version=2
    ```

  - 使用 **Parquet/ORC** 格式,並搭配欄位裁剪 (column pruning) 和述詞下推 (predicate pushdown)。透過壓縮 (compaction) 來管理小檔案,目標是讓每個檔案大小在 128512 MiB 之間,以實現高效掃描。
- Hive metastore
  - 將結構描述和資料表元數據集中在 **Dataproc Metastore** (託管式 Apache Hive Metastore) 或由 **Cloud SQL** 支援的 metastore 中,以便在不同叢集間共享目錄。
  - 使用指向 GCS **外部資料表**以確保耐用性;按日期/小時進行分割,以限制掃描成本。
- 工作、初始化與工作流程
  - 提交 `spark``pyspark``spark-sql`  `hadoop` 工作。**初始化動作 (Initialization actions)** 會在叢集建立時安裝額外的函式庫或代理程式 (例如,連接器、Python 函式庫)
  - **工作流程範本 (Workflow templates)** 可將多步驟的管線參數化;它們可以為每個工作流程建立臨時叢集,然後再將其拆除。這能改善隔離性並降低閒置成本。
  - 建議將**臨時叢集**用於批次 ETL;資料和 metastore 存放在叢集外部 (GCSDataproc MetastoreBigQuery)
- BigQuery 整合
  - **Spark BigQuery 連接器**可直接讀寫 BigQuery;可考慮使用 **BigQuery Storage Read API** 以提高吞吐量,以及使用 **Write API** 以進行低延遲、僅一次 (exactly-once) 的串流插入。
  - 對於資料表維護,可在 BigQuery 中執行下游的 `MERGE`/分割區覆寫操作,以原子化地完成載入。
### Spark 模型、效能調校與可靠性

- API 與執行
  - RDDs:低階、不可變、在 Scala/Java 中具型別安全;您可以控制分割 (partitioning) 與持久化 (persistence)
  - DataFrames/Datasets:關聯式、經 Catalyst 優化;由於其查詢優化與程式碼生成能力,建議優先用於 ETL
  - 轉換 (Transformations) 是延遲執行的 (lazy) ( map, filter, join);動作 (actions) 會觸發執行 ( count, collect, save)Spark 會建立一個由 shuffle 切分的階段 (stages) 所組成的 DAG;任務 (tasks) 會在每個 partition 上執行。
- 分割 (Partitioning)  Shuffle
  - 輸入分割:需要有足夠的 partition 來充分利用所有核心;可從 executor 總核心數的 24 倍開始。透過 `spark.default.parallelism` (用於 RDDs)  reader 選項 (用於 DataFrames) 來控制。
  - Shuffle 分割區:預設值 200 常常會配置過少或過多。調校方式:
--conf spark.sql.shuffle.partitions= {total_executor_cores * 2 to 3}
```
      --conf spark.sql.autoBroadcastJoinThreshold=64m
      ```

    - 為熱點分割區 (hot partitions) 的鍵 (key) 加鹽 (Salt);在 map 端進行預先彙總;及早過濾。
    - 啟用 Adaptive Query Execution (AQE) 來合併 shuffle 後的 partition 並處理傾斜的 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 的容錯模式
  - 冪等寫入 (Idempotent writes):寫入到一個暫存/預備 (staging) 路徑,然後透過目錄層級的 commit 以原子方式提升 (promote);對於 BigQuery,則是寫入預備資料表 (staging table) 然後執行 `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 會將 driver  YARN 的日誌串流至 Cloud Logging;可將其匯出至 sinks 以便進行保留/鑑識。
  - 使用 Cloud Monitoring 指標進行監控:YARN 待處理的 containerCPU、記憶體、HDFS 健康狀況(若有使用)、GCS 吞吐量。針對長時間的 stage 重試、executor 遺失和 speculative execution 突增設定警示。
  - 故障分析:常見原因包括資料傾斜(skew)導致的 stragglershuffle 期間 executor 發生 OOM、物件儲存 commit 失敗,以及 preemptible/spot 節點遺失。謹慎地提高重試次數;過多的重試會增加成本並導致延遲。
- 成本最佳化
  - 使用臨時叢集或 Dataproc Serverless 以避免閒置成本;將資料保存在 GCS 中以最小化永久磁碟的使用。
  - 新增 preemptible/spot 次要 worker 以吸收尖峰需求;設計時應考量到重新計算,因為遺失節點上的任務會被重試。不要將 master 節點放在 preemptible 節點上。
  - 選擇合適的機器類型並使用 autoscaling,在佇列為空時縮減容量。優先使用 Parquet/ORC 搭配分割區裁剪(partition pruning)以降低掃描成本和 CPU 使用率。
  - 透過合併輸出來避免小檔案;較少但較大的檔案可減少 metadata 開銷和工作執行時間。
  - 對於短期的週期性工作(例如,每週 30 分鐘的 Spark ETL),preemptible worker  serverless 通常能提供最佳的成本效益。

#### 實務問題情境

Acme Retail 正在遷移一個 30 節點的本地 Hadoop 叢集,該叢集每晚執行 Spark  Hive ETL,為下游的分析系統提供資料。他們希望以最少的變更來重複使用現有工作、避免全時管理叢集、讓資料的生命週期超越叢集,並降低儲存成本。

方法:
1) 將資料和 metadata 存放在託管服務中
   - 使用 Parquet 格式並搭配分割區(例如,`dt=YYYY-MM-DD`),將所有原始和整理過的資料儲存在 Cloud Storage 中。
   - 理由:GCS 持久、低成本,並將運算與儲存解耦,因此臨時叢集和 serverless 工作可以在沒有永久磁碟的情況下執行。分割區化的 Parquet 可實現謂詞下推(predicate pushdown)和高效掃描。

2) 使用 Dataproc Metastore 集中化目錄
   -  Hive metastore 遷移到 Dataproc Metastore。建立參照 GCS 路徑的外部 Hive 資料表,並保留現有的 schema/分割區邏輯。
   - 理由:託管的 metastore 允許多個臨時叢集和 serverless 工作共享資料表定義,而無需自行運行 HA  MySQL/PostgreSQL 實例。

3) 使用臨時 Dataproc 叢集進行批次 ETL,並使用 workflow templates 進行協調
   - 定義一個 workflow template,它會建立一個具有所需映像檔(例如,`2.1-debian11`)的叢集、執行 Spark 工作(`spark-sql`  `pyspark`),並在完成後刪除叢集。新增初始化動作以安裝任何自訂函式庫。
   - 理由:臨時叢集消除了閒置成本並隔離了工作相依性。Workflow templates 提供了可重複性和參數化(日期、輸入路徑)的能力。

4) 啟用 autoscaling  preemptible worker
   - 附加一個 autoscaling 政策,設定一個小的核心 worker 群組和一個較大的 preemptible 次要 worker 池;調整冷卻時間(cooldowns)以在執行後迅速縮減規模。
   - 理由:核心 worker 維持叢集穩定性;preemptible worker 以較低成本吸收 shuffle 和寬轉換(wide transformations)的負載。Spark/YARN 的重試機制會處理因搶佔而遺失的任務。

5) 透過 Spark BigQuery connector  BigQuery 整合
   - 對於維度/事實載入,將 Spark 結果寫入 BigQuery 的中繼(staging)資料表,然後執行 `MERGE` 陳述式以原子方式更新目標資料表。在可以直接覆寫是安全的情況下,使用分割區覆寫模式(partition overwrite mode)寫入分割區資料表。
   - 理由:BigQuery 可大規模地提供分析和 BI 服務;staging+`MERGE` 的方式可從批次 Spark 實現類似交易的 upsert 操作,從而減少下游的不一致性。

6) 調整 Spark 以提升效能和可靠性
   - 根據 executor 核心數設定 shuffle 分割區並啟用 AQE
 --conf spark.sql.shuffle.partitions=600
 --conf spark.sql.adaptive.enabled=true
 ```
  1. 強化安全性與網路

    • 使用專屬的服務帳號執行叢集,僅授予 GCS 路徑、metastore 和 BigQuery 資料集所需的角色。在受限制的子網路中建立私有 IP 叢集,啟用 Private Google Access,並透過防火牆規則限制 UI 存取。
    • 理由:最低權限和網路隔離可減少攻擊面;私有的控制平面出口可避免公開暴露。
  2. 建置日誌、歷史記錄與警示

    • 啟用 Spark 事件日誌到 GCS 並部署 History Server;將 driver/YARN 日誌路由到 Cloud Logging 並設定保留政策。新增 Monitoring 警示,用於監控長時間等待的 container、重複的任務失敗或過長的工作執行時間。
    • 理由:集中化的日誌支援根本原因分析;主動式警示能及早偵測到資料傾斜、OOM 或 I/O 效能下降等問題。
  3. 針對臨時性需求和彈性尖峰,選擇性地使用 Dataproc Serverless 進行現代化

    • 將零星或探索性的 Spark SQL 工作負載移至 Dataproc Serverless;在 serverless 上完全驗證之前,將夜間的 pipeline 維持在臨時叢集上執行。
    • 理由:Serverless 免除了叢集維運並能自動擴展,非常適合不可預測的負載;現有工作流程只需極少的程式碼變更即可繼續運作。
  4. 驗證物件儲存 committer 和小檔案管理

    • 設定 FileOutputCommitter algorithm v2,並在寫入前透過 repartition/coalesce 將輸出合併為每個檔案 256–512 MiB。
    • 理由:物件儲存缺乏原子性的 rename 操作;最佳化的 committer 可減少複製/重新命名的開銷。檔案合併可緩解小檔案問題,從而提升效能並降低成本。

此設計以最少的重構重複利用了現有的 Spark 和 Hive 工作,確保了 GCS 中資料的持久性,集中化了 schema,控制了安全性的爆炸半徑,提供了強大的可觀測性,並透過臨時叢集、autoscaling、preemptible 容量以及針對性地使用 serverless 執行來最佳化成本。


訊息傳遞、事件擷取與即時服務 · 所有領域 · 資料擷取、整合與遷移

練習這些題目 → · 在 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

一站式存取

一份訂閱。所有考試。

每個方案都可無限存取答案搜尋、練習測驗、AI 解釋和完整的資源庫 — 支援 20 多種語言。

每月
24.87
Just €0.83/day
包含所有內容:
  • 無限答案搜尋
  • 無限練習測驗
  • AI 驅動的解釋
  • 完整資源庫
  • 20 多種語言
  • 每週內容更新
  • 獎勵與推薦
  • 優先支援
開始免費試用

無需信用卡*

最佳價值
12 個月
179.87
Just €0.49/daySave 40%
包含所有內容:
  • 無限答案搜尋
  • 無限練習測驗
  • AI 驅動的解釋
  • 完整資源庫
  • 20 多種語言
  • 每週內容更新
  • 獎勵與推薦
  • 優先支援
開始免費試用

無需信用卡*

✓ 包含免費方案 · ✓ 隨時取消 · ✓ 所有方案解鎖完整產品