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
- 叢集類型與節點角色
- 主要 (master) 節點託管 YARN、HDFS NameNode (若有使用) 及 Spark 驅動程式的 UI;HA 模式會使用多個主要節點。
- 工作節點執行 executor 與 HDFS DataNode (若有使用)。
- 次要/輔助工作節點通常是可搶佔/Spot 執行個體,用於提供具彈性、成本較低的容量,且不具備 HDFS 角色。
- 映像檔綑綁了作業系統和元件版本 (例如 2.1-debian11、2.2-ubuntu20);鎖定映像檔版本以控制 Spark/Hadoop 的相容性,並謹慎地進行升級。
- Component Gateway 透過 HTTPS 安全地發布 UI (Spark History Server、YARN RM)。
- Dataproc Serverless for Spark
- 無需佈建叢集,具備自動擴縮功能,並對 executor 和 driver 採秒級計費。非常適合零星或突發性的工作,或是在需要將維運負擔降到最低時使用。
- 權衡之處:可供調整的底層選項比叢集少;工作啟動延遲可能高於已暖機的叢集;使用 serverless 指標和事件日誌進行故障排除。
- 自動擴縮
- 叢集自動擴縮政策會根據 YARN/Spark 指標和冷卻時間來新增/移除工作節點,並可分別調整主要和次要工作節點群組。
- Serverless 的自動擴縮由服務本身管理;設計時應使其具備分割區平行處理能力,並避免序列化的瓶頸,以達到最佳的擴縮效果。
- 儲存與連接器
- 建議使用 Google Cloud Storage (GCS) 作為記錄系統 (system-of-record);它將運算與儲存解耦、降低永久磁碟成本,且能在叢集生命週期結束後繼續存在。
- GCS 連接器 (gs://) 與 Hadoop/Spark 整合。寫入物件儲存區時會使用提交協定;設定 FileOutputCommitter 演算法為 v2,以減少重新命名的負擔並加速在 GCS 上的工作提交:
--conf mapreduce.fileoutputcommitter.algorithm.version=2
```
- 使用 **Parquet/ORC** 格式,並搭配欄位裁剪 (column pruning) 和述詞下推 (predicate pushdown)。透過壓縮 (compaction) 來管理小檔案,目標是讓每個檔案大小在 128–512 MiB 之間,以實現高效掃描。
- Hive metastore
- 將結構描述和資料表元數據集中在 **Dataproc Metastore** (託管式 Apache Hive Metastore) 或由 **Cloud SQL** 支援的 metastore 中,以便在不同叢集間共享目錄。
- 使用指向 GCS 的**外部資料表**以確保耐用性;按日期/小時進行分割,以限制掃描成本。
- 工作、初始化與工作流程
- 提交 `spark`、`pyspark`、`spark-sql` 或 `hadoop` 工作。**初始化動作 (Initialization actions)** 會在叢集建立時安裝額外的函式庫或代理程式 (例如,連接器、Python 函式庫)。
- **工作流程範本 (Workflow templates)** 可將多步驟的管線參數化;它們可以為每個工作流程建立臨時叢集,然後再將其拆除。這能改善隔離性並降低閒置成本。
- 建議將**臨時叢集**用於批次 ETL;資料和 metastore 存放在叢集外部 (GCS、Dataproc Metastore、BigQuery)。
- 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 總核心數的 2–4 倍開始。透過 `spark.default.parallelism` (用於 RDDs) 和 reader 選項 (用於 DataFrames) 來控制。
- Shuffle 分割區:預設值 200 常常會配置過少或過多。調校方式:
--conf spark.sql.shuffle.partitions= {total_executor_cores * 2 to 3}
```
- 目標是在寬轉換 (wide transforms) 後,每個 partition 大小約為 100–256 MiB;太小會導致排程器額外負擔 (overhead),太大則有 executor 發生記憶體不足 (OOM) 的風險。
- Shuffle 是 joins、groupBy 和 orderBy 的主要成本。確保 executor 有足夠的記憶體和磁碟空間;在叢集上針對繁重的 shuffle 操作可考慮使用本地 SSD。
- 傾斜 (Skew) 與 Join 策略
- 偵測資料傾斜 (任務執行時間出現長尾現象、partition 大小過大)。緩解方法:
- 廣播 (Broadcast) 小表以避免 shuffle:
- 偵測資料傾斜 (任務執行時間出現長尾現象、partition 大小過大)。緩解方法:
--conf spark.sql.autoBroadcastJoinThreshold=64m
```
- 為熱點分割區 (hot partitions) 的鍵 (key) 加鹽 (Salt);在 map 端進行預先彙總;及早過濾。
- 啟用 Adaptive Query Execution (AQE) 來合併 shuffle 後的 partition 並處理傾斜的 join:
--conf spark.sql.adaptive.enabled=true
```
- 快取 (Caching)、檢查點 (Checkpointing) 與血緣 (Lineage)
- 當熱點的中繼 DataFrame 會被重複使用時,才謹慎地快取它;優先使用
MEMORY_AND_DISK以避免 OOM。 - 將長的血緣關係設置檢查點到 GCS 或 HDFS,以限制失敗時的重新計算範圍。
- 當熱點的中繼 DataFrame 會被重複使用時,才謹慎地快取它;優先使用
- Executor 與動態配置
- 適當調整 Executor 的大小,以平衡平行度與 GC 的額外負擔:
- 每個 Executor 的核心數:對於平衡 I/O 與 CPU 的任務,設定為 2–5 個核心;較少的核心數可減少 GC 暫停時間。
- 記憶體額外開銷:對於寬 shuffle,設定
spark.yarn.executor.memoryOverhead。 - 在叢集上啟用外部 shuffle 服務並搭配動態配置,讓 executor 能根據工作負載進行擴展:
- 適當調整 Executor 的大小,以平衡平行度與 GC 的額外負擔:
--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 (...)
```
- 增量處理:在
ingestion_date分割區上使用基於浮水印 (watermark) 的過濾;在 GCS 中維護一個已處理清單 (processed-manifest) 以避免重複處理。 - 無效信件 (Dead-letter) 處理:當發生解析/驗證錯誤時,將錯誤的記錄連同診斷資訊分流到一個隔離路徑/資料表。若要嚴格強制 schema 並使用內建的 DLQ,可考慮使用 Dataflow;若使用 Spark,則需實作針對每筆記錄的 try/catch 並導向一個獨立的 sink。
安全性、可觀測性與成本
- 身分與存取權限
- 在專屬的服務帳號下執行叢集與工作,並遵循最低權限 IAM 原則。僅授予必要的角色,例如:
roles/dataproc.worker給予執行個體服務帳號roles/storage.objectViewer或objectAdmin給予 GCS I/O 路徑roles/bigquery.dataEditor給予目標資料集
- 對於 Dataproc Serverless,使用每個工作專屬的服務帳號來限定存取範圍。
- 在專屬的服務帳號下執行叢集與工作,並遵循最低權限 IAM 原則。僅授予必要的角色,例如:
- 網路隔離與加密
- 在 VPC 子網路中使用私有 IP 叢集,透過防火牆限制 master UI 的存取,並啟用 Private Google Access 以便在沒有公開出口的情況下存取 GCS/BigQuery。
- 將叢集放置在 Shared VPC 專案中以進行集中控管。可選擇在 Dataproc 上啟用 Kerberos 以進行叢集內驗證。
- 使用 CMEK 進行靜態加密:在 GCS 儲存桶、永久磁碟、Dataproc Metastore 和 BigQuery 上設定 CMEK;傳輸中預設使用 TLS 加密。
- 日誌、歷史記錄與指標
- 啟用 Spark 事件日誌到 GCS 並部署 History Server:
--conf spark.eventLog.enabled=true
--conf spark.eventLog.dir=gs://bucket/spark-events/
```
- Dataproc 會將 driver 和 YARN 的日誌串流至 Cloud Logging;可將其匯出至 sinks 以便進行保留/鑑識。
- 使用 Cloud Monitoring 指標進行監控:YARN 待處理的 container、CPU、記憶體、HDFS 健康狀況(若有使用)、GCS 吞吐量。針對長時間的 stage 重試、executor 遺失和 speculative execution 突增設定警示。
- 故障分析:常見原因包括資料傾斜(skew)導致的 straggler、shuffle 期間 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
```
- 對小型維度表使用廣播聯結(broadcast join),並將長的 lineage 檢查點(checkpoint)存到 GCS 以提高穩定性。
- 理由:適當的分割區可減少資料傾斜和排程器開銷;AQE 在執行期會根據資料特性進行調整;檢查點機制可限制故障後重新計算的範圍。
強化安全性與網路
- 使用專屬的服務帳號執行叢集,僅授予 GCS 路徑、metastore 和 BigQuery 資料集所需的角色。在受限制的子網路中建立私有 IP 叢集,啟用 Private Google Access,並透過防火牆規則限制 UI 存取。
- 理由:最低權限和網路隔離可減少攻擊面;私有的控制平面出口可避免公開暴露。
建置日誌、歷史記錄與警示
- 啟用 Spark 事件日誌到 GCS 並部署 History Server;將 driver/YARN 日誌路由到 Cloud Logging 並設定保留政策。新增 Monitoring 警示,用於監控長時間等待的 container、重複的任務失敗或過長的工作執行時間。
- 理由:集中化的日誌支援根本原因分析;主動式警示能及早偵測到資料傾斜、OOM 或 I/O 效能下降等問題。
針對臨時性需求和彈性尖峰,選擇性地使用 Dataproc Serverless 進行現代化
- 將零星或探索性的 Spark SQL 工作負載移至 Dataproc Serverless;在 serverless 上完全驗證之前,將夜間的 pipeline 維持在臨時叢集上執行。
- 理由:Serverless 免除了叢集維運並能自動擴展,非常適合不可預測的負載;現有工作流程只需極少的程式碼變更即可繼續運作。
驗證物件儲存 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.
通過考試 →