Google PDE: 資料擷取、整合與遷移 — 學習指南
屬於 Google Professional Data Engineer — 學習指南. 使用經過驗證的解答練習: Google 考試中心, 或參加限時模擬考試: ExamRoll.io.
總覽
在 Google Cloud 中,資料擷取、整合與遷移涵蓋了可重複的模式、託管服務與操作控制,能將多樣的來源系統轉化為可靠、可查詢的資料集。有效的設計會將傳輸與轉換分離、將生產者與消費者解耦,並偏好使用具備冪等性、檢查點、清晰血緣關係與驗證機制的管線。本節將涵蓋擷取模式、用於資料移動和 CDC 的 Google Cloud 服務、結構描述與資料品質控制、連線能力與混合雲整合,以及轉換策略,並在全文中點出設計上的權衡取捨與失敗模式。
擷取模式與工作負載
- 批次擷取 (Batch ingestion):在定義的時間間隔內定期拉取或放置檔案。適合可預測的成本與資料回填。失敗模式:大型、不頻繁的批次會導致資源使用量遽增、追趕資料的空窗期過長,以及錯過 SLA。緩解措施:適當調整批次大小、依時間或鍵值進行分片,並使用平行處理。
- 大量載入 (Bulk load):一次性或大規模的載入(例如,初期的歷史資料回填)。偏好使用欄位式或自描述格式(Parquet、Avro),並直接載入到分析型儲存(BigQuery)或暫存於 Cloud Storage。權衡取捨:查詢外部資料表可避免載入步驟,但會將成本轉移到查詢時的掃描。
- 增量載入 (Incremental load):透過時間戳記或高水位標記 (high-water marks) 定期進行差異載入。需要穩健的重複資料刪除與冪等的 upsert(更新或插入)操作。失敗模式:時鐘偏移或延遲到達的記錄。應使用伺服器端的提交時間戳記與浮水印 (watermarking)。
- 異動資料擷取 (Change data capture, CDC):從營運資料庫中持續複製插入、更新與刪除操作。最適合近乎即時的分析與低停機時間的遷移。權衡取捨:
- 順序性:大多數 CDC 工具會保留交易內的順序,且通常在單一分片內也能維持順序,但不保證跨分片的全域順序。應使用交易提交時間戳記與主鍵來重構順序。
- 交付語意:通常為「至少一次」(at-least-once);應建構冪等的接收端 (sink) 或使用唯一的變更 ID 來進行重複資料刪除。
- 快照 + CDC:從一個一致性的快照開始,然後從精確的日誌序列套用變更,以在不停機的情況下達到資料同步。
關聯式、SaaS、地端與檔案來源:
- 關聯式資料來源:使用原生的 CDC 或時間戳記欄位。對於大量載入,可匯出為 Avro/Parquet 格式並暫存於 Cloud Storage。
- SaaS 來源:偏好使用帶有增量 token 的供應商 API;透過託管連接器(例如,在 Data Fusion 中)進行整合。針對速率限制進行節流,並處理結構描述的變動。
- 地端來源:可選擇基於代理程式的傳輸、VPN/Interconnect + Private Google Access,或使用 Transfer Appliance 進行離線植入。
- 檔案擷取:對於大量小檔案,可將其打包(例如,使用 tar)以減少 RPC 的開銷。使用
gsutil -m或平行化用戶端;將檔案組合或轉換為較大的欄位式檔案以利分析。
用於擷取、整合與遷移的 Google Cloud 服務
- Datastream (無伺服器 CDC):從 MySQL、PostgreSQL 和 Oracle 擷取變更到 Cloud Storage、BigQuery (透過範本) 或 Pub/Sub。它會保留交易邊界與提交的中繼資料;不保證全域順序。應在下游根據鍵值與提交時間戳記來排序。預期為「至少一次」交付;應設計冪等的消費者(例如,使用 BigQuery 的
MERGE搭配變更 ID)。 - Database Migration Service (DMS):用於使用原生複製技術進行停機時間極短的資料庫遷移。DMS 會建立一個一致性的快照,然後使用 GTID/LSN/SCN 持續複製變更。它是專為直接遷移 (lift-and-shift) 而設計,不適用於任意的轉換。若用於分析,可視需求搭配 Dataflow 或 Data Fusion 來擴充 DMS 的功能。
- Cloud Data Fusion:一個託管的整合服務,提供連接器以連接關聯式資料庫、SaaS、檔案與訊息傳遞系統。可建構包含轉換階段(連接、彙總、格式轉換、自訂 Wrangler 配方)的管線,並擷取跨來源與欄位的血緣關係。在操作上,它會進行排程、重試並發送指標。可使用 Data Fusion 進行無程式碼/低程式碼的 ELT/ETL,並集中管理連接器。
- Storage Transfer Service (STS):從 AWS S3、Azure Blob、地端(使用代理程式)、SFTP 與 URL 列表,將資料進行託管、排程傳輸至 Cloud Storage。支援資訊清單 (manifests)、增量同步、頻寬控制與校驗和完整性檢查。失敗模式包含小檔案效率不彰與 API 節流;可透過批次處理與可調整的並行性來緩解。
- Transfer Appliance:一種離線、加密的設備,用於在網路頻寬有限或資料過於敏感不宜長時間傳輸時,進行 TB 到 PB 等級的初始資料植入。內建監管鏈 (Chain-of-custody) 與加密機制。植入資料後,可接著使用 STS 或 CDC 來處理差異資料。
- Cloud Pub/Sub + Dataflow:Pub/Sub 將生產者與消費者解耦,適用於串流或微批次模式。Dataflow 提供自動擴展、具備狀態的串流/批次處理,並帶有檢查點與浮水印功能。使用 BigQuery Storage Write API 進行低延遲串流,可在預設串流中提供「恰好一次」(exactly-once) 的保證;否則需依賴
insertId的重複資料刪除語意。
對於 Hadoop 到 Dataproc 的遷移,應透過 GCS 連接器將資料儲存在 Cloud Storage 中,並使用臨時性或自動擴展的叢集,以最小化 Persistent Disk 的使用。這樣可以避免高昂的區塊儲存成本,同時保留與 HDFS 相容的處理語意。
邊界上的結構描述、驗證與資料品質
- 結構描述對應與類型轉換:及早標準化為強型別結構描述。Avro 或 Parquet 能保存結構描述並乾淨地演進。在 BigQuery 中,優先使用分區和叢集資料表以降低掃描成本。 範例:為每日分析建立一個分區資料表
undefined
- 格式錯誤記錄的處理:將拒絕的記錄路由到 dead-letter queue (Pub/Sub) 或 Cloud Storage 中的隔離儲存桶。在 Dataflow 中使用側邊輸出,或在 Data Fusion 中使用錯誤收集器。記錄剖析錯誤,並附上範例 payload 和結構描述版本以便分類處理。
- 驗證:在持久化之前執行邊界檢查:
- 結構性:結構描述一致性、必要欄位、資料類型、列舉範圍。
- 參考性:透過快取維度查詢來確認外鍵是否存在。
- 合理性:時間戳記的範圍、地理圍欄、非負金額。
- 唯一性:主鍵或複合鍵衝突。
- 冪等載入:使用確定性鍵和 upsert 操作。在 BigQuery 中,使用自然鍵或代理變更鍵來實作 MERGE。 範例:
undefined
- 浮水印與延遲:在串流管線中,設定事件時間浮水印和允許的延遲,以平衡完整性和延遲。延遲的資料會被路由到修正路徑或觸發回填作業。
- 對帳:從來源到接收端,追蹤每個分區/視窗的資料列計數和校驗和。擷取 CDC 日誌位置 (LSN/SCN) 和提交時間戳記;儲存在控制表中以證明連續性並識別間隙。
連線能力、可靠性與維運
網路連線與私有存取:
- 混合雲:使用 Cloud VPN 或 Dedicated/Partner Interconnect 進行私有連線。啟用 Private Google Access 或 Private Service Connect 以私有方式存取 Google API (如 Cloud Storage)。
- 安全性:使用服務帳戶作為工作負載身分、最小權限 IAM、VPC Service Controls 以防止資料外洩,並在需要時使用 CMEK。
- 吞吐量:在用戶端擴展平行處理能力,但最終由頻寬決定吞吐量。對於大規模傳輸,初期的大量資料建議使用 Transfer Appliance,然後使用 STS 或 CDC 進行增量更新。
檢查點與反壓:Dataflow 會管理檢查點和自動擴展;設計能夠吸收突發流量的接收端 (例如緩衝到 Cloud Storage、批次寫入 BigQuery)。對於 Pub/Sub,調整流量控制和確認期限,以防止訊息重傳風暴。
使用 CDC 的排序與一致性:
- Datastream 會保留交易內的順序並發出提交中繼資料;消費者使用提交時間戳記來重構每個鍵的順序。預期為 at-least-once (至少一次);需建立冪等性。
- DMS 使用原生誌在快照和複寫切換之間確保資料庫一致性。使用讀取複本或雙寫策略進行分階段切換。
分析用的檔案策略:對於大型、多引擎存取,將標準資料儲存在 Cloud Storage 中,並在符合成本效益的情況下,公開永久性外部資料表以供臨時查詢。對於生產分析,則載入到 BigQuery 的分區資料表中,以最小化每次查詢的掃描成本。
小檔案優化:在傳輸前將小檔案打包 (例如,每個 tar 檔約 1,000 個),然後在雲端解壓縮。使用平行的 gsutil 和生命週期規則來分層和過期暫存成品。
維運陷阱與緩解措施:
- 來自 SaaS 的結構描述漂移:在 Data Fusion 中啟用結構描述演進並強制執行相容性。對破壞性變更發出警示。
- 時區與編碼:在入口處標準化為 UTC 和 UTF-8。
- CDC 中的間隙:監控來源日誌的保留期;當複本延遲接近保留限制時發出警示。
- 配額:BigQuery 串流插入、API 速率限制;當接近限制時進行批次處理。
切換 (Cutover)、回填 (Backfill) 與驗證 (Verification)
- 切換計畫:
- 大爆炸式 (Big bang):短暫凍結、一次性切換。營運複雜度最低;但若需回復 (rollback) 風險最高。
- 分階段或藍綠部署 (Phased or blue/green):雙軌運行,搭配鏡像寫入、漸進式流量轉移及影子讀取 (shadow reads)。成本較高;但回復 (rollback) 較安全。
- 回填 (Backfill):
- 執行初始的大量載入 (bulk load)(使用 Transfer Appliance 或 STS),並採用 Avro/Parquet 格式以保留 schema。在載入期間進行分割 (partition) 與叢集 (cluster),以避免後續重工。
- 在快照的同時,從一個已知的日誌位置啟動 CDC,以擷取大量傳輸期間的差異資料 (delta)。在對生產環境開放前,於一個共同的浮水印 (watermark) 點進行核對。
- 回復 (Rollback):
- 在驗證期間,將舊有系統維持唯讀狀態。對於雙寫入情境,將寫入操作置於功能旗標 (feature flag) 後方,以便快速還原。保留一個一致的檢查點 (checkpoint),以便在需要時重播 (replay) 或撤銷 (unwind) CDC 的變更。
- 遷移驗證:
- 結構性:資料列計數與每個分割區的校驗和 (checksum) 相符;schema 與約束條件 (constraints) 對等。
- 時間性:從快照邊界到切換點無資料間隙;CDC 位置是連續的。
- 業務對等性:比較不同時間視窗內的匯總資料與 KPI;執行驗收查詢 (acceptance queries)。
- 效能:根據預算驗證資料擷取吞吐量、查詢延遲與成本。
實務問題情境
Northstar Retail 公司必須將全球混合的本地端 Oracle 和 MySQL 交易系統、SaaS CRM 事件以及每日的 CSV 檔案,整合到 Google Cloud 中,以支援近乎即時的分析與機器學習。他們還需要遷移一個舊有的 Hadoop 叢集,同時避免產生高昂的區塊儲存費用,並實現零或低停機時間的切換。
- 建立安全的混合式連線
- 使用 Partner Interconnect 作為主要頻寬,並以 Cloud VPN 作為備援。啟用 Private Google Access,讓本地端工作負載能以私密方式存取 Cloud Storage 與 Pub/Sub。 理由:私密路徑能將出口 (egress) 的曝險與延遲降至最低,而 Private Google Access 則無需公開 IP 即可滿足安全政策要求。
- 有效率地植入歷史資料
- 對於 800 TB 的歷史 HDFS 資料,使用 Transfer Appliance 將其複製到 Cloud Storage(初始大量載入)。植入資料後,每日從本地端的 NFS 匯出執行 Storage Transfer Service,以擷取變更,直到切換為止。 理由:Transfer Appliance 避免了長時間的網路飽和;STS 提供排程化、具校驗和的增量同步。將資料儲存在 Cloud Storage 並搭配 GCS connector,讓 Dataproc 無需為每個節點配置 50 TB 的 Persistent Disk 即可進行處理。
- 透過 CDC 遷移營運資料庫
- 使用 DMS 以最低停機時間遷移 MySQL 與 PostgreSQL。對於 Oracle 到分析系統的 CDC,使用 Datastream 將資料落地到 Cloud Storage,然後透過 Google 提供的 Dataflow 範本載入到 BigQuery。 理由:DMS 利用原生複製功能,提供可靠的快照 + 持續同步;Datastream 提供無伺服器的 CDC 並帶有提交中繼資料 (commit metadata),而 Dataflow 範本確保寫入 BigQuery 的操作是有序且冪等的 (idempotent)。
- 擷取 SaaS 與檔案型饋送
- 建立 Cloud Data Fusion 管線,使用 SaaS 連接器搭配增量權杖 (incremental tokens) 來處理 CRM 事件,並建立一個檔案管線,透過 STS 從供應商的 SFTP 伺服器擷取每日的 CSV 檔案。將資料正規化為 Avro 格式,存放在一個策劃過的 (curated) Cloud Storage 儲存桶中,然後載入到分割的 BigQuery 資料表。 理由:Data Fusion 將連接器、轉換與資料血緣 (lineage) 集中管理。標準化為 Avro 格式可保留 schema 並簡化後續演進;分割的 BigQuery 資料表可降低查詢成本。
- 串流即時事件
- 將網站與商店事件發布到 Pub/Sub。使用 Dataflow 進行解析、驗證、擴充 (enrichment) 與浮水印處理;透過 Storage Write API 寫入 BigQuery,並將原始的 Avro 檔案封存至 Cloud Storage。 理由:Pub/Sub 將生產者與消費者解耦;Dataflow 提供自動擴展、有狀態處理、檢查點與延遲資料處理;雙寫入確保了低延遲分析與持久的原始資料保留。
- 強制執行邊界資料品質與 schema 控制
- 在 Dataflow/Data Fusion 中實作 schema 註冊與驗證。將格式錯誤的記錄路由到一個 GCS 的隔離儲存桶與 Pub/Sub 的無效信件主題 (dead-letter topic)。應用領域檢查(例如:貨幣代碼、UTC 時間戳)並使用複合鍵 (composite keys) 進行重複資料刪除。 理由:及早拒絕與隔離可防止壞資料擴散;冪等性 (idempotency) 與重複資料刪除可防範來自 CDC 與串流來源的「至少一次」傳遞 (at-least-once delivery) 問題。
- 優化分析儲存與存取
- 將策劃過的資料集載入到經過分割與叢集處理的 BigQuery 資料表中。將原始封存資料公開為永久的外部資料表,供低頻率的探索性查詢使用。對於仍需保持交易性的 OLTP 工作負載,保留 Cloud SQL 並搭配讀取複本 (read replicas)。 理由:分割與叢集能將掃描成本降至最低;外部資料表避免了為偶爾存取而進行不必要的載入;Cloud SQL 為交易型應用程式保留了 ACID 語意。
- 規劃切換、回填與回復
- 為每個 RDBMS 執行快照 + CDC;達到一個核對點,確保資料列計數與校驗和相符。執行藍綠部署並進行 48 小時的雙寫入,同時逐步將讀取流量轉移到 BigQuery。維持一個功能旗標,以便在偵測到差異時還原寫入操作。 理由:藍綠部署降低了風險;在已知的浮水印點進行驗證確保了完整性;功能旗標讓快速回復成為可能。
- 驗證與可觀測性
- 建立控制表,擷取每個分割區的來源 LSN/SCN、提交時間戳、資料列計數與校驗和。監控 Datastream 延遲、DMS 複製狀態、Dataflow 浮水印、Pub/Sub 積壓量 (backlog)、STS 工作狀態以及 BigQuery 串流插入指標。 理由:端到端的資料血緣與量化控制提供了可稽核的正確性證明,並能針對資料間隙或延遲及時發出警報。
透過分離落地 (landing)、策劃 (curation) 與服務 (serving) 層;使用 Cloud Storage 作為耐用、低成本的中繼與封存儲存;利用 DMS/Datastream 進行 CDC 並搭配冪等的消費者;以及在入口處強制執行 schema 與品質控管,Northstar Retail 得以實現安全、可擴展的資料擷取,以及一個低風險、可驗證且成本可預測的遷移。
← Spark、Dataproc 與分散式資料處理 · 所有領域 · 工作流程編排與管線自動化 →
練習這些題目 → · 在 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.
通過考試 →