Google PDE: 工作流程編排與管線自動化 — 學習指南
屬於 Google Professional Data Engineer — 學習指南. 使用經過驗證的解答練習: Google 考試中心, 或參加限時模擬考試: ExamRoll.io.
可靠性、故障處理與冪等性
重試、逾時與退避 (backoff)
- 對於暫時性故障,請使用有界指數退避 (bounded exponential backoff),並將總重試時間窗口限制在作業的 SLA 內。例如,一個每 15 分鐘輪詢一次資料庫的前端或任務,應該使用指數退避重試最多 15 分鐘,然後呈現一個受控的失敗。
- 在 Airflow 中,設定每個任務的
execution_timeout和全域 DAG 的 SLA;在 Workflows 中,設定每個步驟的逾時和重試策略,並使用max_doublings和max_retry_duration。對於 Cloud Run jobs,設定重試次數和退避。
回填 (Backfills)、追補 (Catchup) 與故障處理
- 當任務具備冪等性且來源資料已按日期分區時,啟用追補 (catchup) 功能以重新計算歷史資料。對於非確定性輸出或有外部副作用的情況,請考慮使用僅供回填的 DAG 或寫入稽核表來追蹤已產生的內容。
- 對於串流/批次轉換中的記錄層級故障,請使用死信主題/資料表 (dead-letter topics/tables)。對於批次 Dataflow,使用錯誤標籤捕獲不良資料列並匯總錯誤指標;對於串流處理,則使用 Pub/Sub 的 DLQ (Dead-Letter Queue)。
冪等性任務設計與重新執行
- BigQuery:偏好使用
MERGE或帶有去重複鍵的INSERT;使用insertId來為串流插入進行去重複。對於批次處理,先寫入一個暫存表,然後在一個具備交易安全性的步驟中MERGE到目標表,以允許完整的重新執行。 - Cloud Storage:使用世代前置條件 (generation preconditions) 和確定性的物件名稱 (例如:
prefix/date/hash),這樣重新執行時只會在預期情況下安全地覆寫。 - Pub/Sub 和 Dataflow:設計時需考慮至少一次 (at-least-once) 的交付。在訊息中包含識別碼 (例如,包裹 ID、邏輯事件時間戳),以便下游系統可以去重複並處理延遲問題。如果業務規則接受「先處理的事件獲勝」的語意,請記錄此權衡並監控資料傾斜;否則,應根據事件時間並搭配決勝規則 (tie-breakers) 來決定獲勝者。
- 從部分失敗中恢復:按
run_id或日期對輸出進行分區,寫入完成標記,並讓下游任務依賴這些標記。僅重新處理標記為未完成的分區。
故障排除與擴展性
- 當串流儀表板遺漏事件,但 Pub/Sub 中顯示事件存在時,可以將一組已知的固定資料集送入 Dataflow pipeline,以隔離轉換邏輯中的缺陷。請驗證視窗 (windowing)、觸發器 (triggers) 和允許的延遲 (allowed lateness) 設定。
- 常見的失敗模式:為無邊界資料源建立串流 pipeline 時,若沒有設定適當的視窗/觸發器,或不正確地使用分片視窗 (sharded window),可能導致 pipeline 建立失敗或狀態 (state) 爆炸。
- 透過
max workers和自動擴展演算法來擴展 Dataflow;對於流量尖峰 (例如,50,000 次安裝),應提高最大工作節點數,以允許在高峰期間進行水平擴展。
安全性、參數化、環境與 CI/CD
參數化與組態管理
- 將組態依環境外部化。在 Composer 中,使用 Variables、Connections 和環境變數;根據執行日期或分區來模板化 DAG 參數。在 Workflows 中,使用執行時參數 (runtime arguments),並為每個環境使用獨立的 workflow,或從 Secret Manager 讀取組態。
- 透過讀取一個控制表 (例如,列出客戶、來源或分區的 BigQuery 組態資料集) 來實現元資料驅動的協調。動態生成任務,使程式碼變更與資料驅動的變更解耦。
密鑰、服務帳戶與最小權限原則
- 將憑證儲存在 Secret Manager 中,並在執行時引用。避免將密鑰寫死在程式碼或 Airflow Variables 中。
- 為每個 pipeline 指派一個獨立的服務帳戶,並賦予其所需的最小 IAM 角色。對於受監管的 BigQuery 存取,應將客戶資料隔離到不同的資料集中,僅將特定於資料集的角色授予經批准的使用者,並限制只有經批准的 principals 才能存取 BigQuery API。對於多租戶架構,為每個客戶建立一個資料集,並僅綁定適當的角色。
CI/CD 與基礎設施即程式碼 (IaC)
- 使用 Terraform 管理基礎設施 (Composer 環境、Workflows、Scheduler 作業、Pub/Sub 主題、日誌接收器)。使用模組 (modules) 來標準化專案/環境、密鑰和服務帳戶的設定。
- 使用 Cloud Build 或 GitHub Actions 來建置和測試 pipeline 程式碼。自動化單元測試、SQL 語法檢查 (linting)、Dataform 試運行 (dry-runs) 以及 Airflow DAG 驗證。透過標籤 (tags) 來晉升交付產物 (artifacts);對於 Composer,將 DAG 打包成可部署的 bundles;對於 Dataform,使用在斷言 (assertions) 通過後才晉升的發布分支。
- 部署晉升:透過獨立的專案和參數化的組態,實現 dev → test → prod 的流程。對於高風險的晉升,採用持續交付 (continuous delivery) 搭配手動批准關卡和變更窗口。
可觀測性、警報與 Runbook
遙測與警報
- 將所有 orchestration 的日誌以結構化欄位 (pipeline, dag_id, run_id, task_id, partition) 傳送到 Cloud Logging。透過以日誌為基礎的指標 (log‑based metrics),將錯誤日誌匯出到 Monitoring。並針對以下情況設定警報:
- 錯過排程或未達成 SLA
- 連續任務失敗
- 待辦項目增長 (例如:Pub/Sub 未確認訊息、Dataflow 系統延遲)
- 資料品質斷言失敗
- Cloud Composer:監控 DAG/任務的持續時間、成功率、佇列深度與 scheduler 健康狀況。設定 on_failure_callback 以觸發呼叫通知 (paging) 與修復用的 runbook。
- Cloud Workflows:檢查執行日誌 (Execution logs) 與步驟延遲;加入明確的重試與錯誤處理機制;發送帶有關聯 ID (correlation ID) 的自訂日誌。
- BigQuery 資料表變更通知:建立一個專案層級的 Logging sink,並使用進階篩選器來鎖定針對特定資料表的插入作業 (insert jobs),然後匯出到 Pub/Sub;您的監控工具訂閱該主題,即可獲得即時警報,且不受其他資料表雜訊的干擾。
Runbook 設計
- 為每個 pipeline 文件化其觸發器、相依性、SLA、回滾/重試程序,以及安全的回填 (backfill) 步驟。內容應包含 Dataflow 的「固定資料集重播 (fixed dataset replay)」、如何排空 (drain) 串流作業、如何重新處理失敗的分割區 (partition),以及如何修復 DLQ (Dead-Letter Queue) 中的訊息。
- 透過決策樹與上報路徑 (escalation paths),捕捉常見的失敗特徵 (例如:權限遭拒、配額超限、schema 不符)。
實務問題情境
Acme 零售分析公司需要擷取每日合作夥伴提供的 CSV 檔案,這些檔案偶爾會包含格式錯誤的資料列。他們需要轉換並將有效的資料載入到 BigQuery,並將錯誤的資料列呈現出來以供調查。此外,他們還希望透過事件驅動的方式進行資料擴充,以實現近乎即時的價格更新,並能安全地將程式碼從開發 (dev) 環境推廣到生產 (prod) 環境。
方法:
儲存與事件觸發器
- 建立一個專用的 Cloud Storage 儲存貯體,並啟用物件版本管理 (object versioning) 與統一的儲存貯體層級存取權 (uniform bucket‑level access)。透過 Eventarc 將物件完成 (object finalize) 的通知發送到 Pub/Sub。
- 理由:物件完成 (Object finalization) 是一個可靠的事件,可用於觸發下游的擷取作業;版本管理則支援重新執行與稽核。
具備無效信件處理 (dead‑letter handling) 的批次擷取
- 使用 Cloud Composer 來排程一個每日 02:00 執行的 Airflow DAG,並啟用 catchup。此 DAG 會啟動一個 Dataflow 批次作業,該作業會解析 CSV、驗證 schema,並使用確定性的中繼暫存資料表 (deterministic staging tables) 將有效記錄寫入 BigQuery,然後透過 MERGE 指令合併到分割後的目標資料表中。將格式錯誤或失敗的記錄路由到一個 BigQuery 的無效信件資料表 (dead-letter table)。
- 理由:Dataflow 可擴展解析/驗證的效能;MERGE 確保了冪等性 (idempotency);無效信件的捕獲支援在不阻斷 pipeline 的情況下進行檢查,這符合處理格式錯誤資料列的建議模式。
事件驅動的資料擴充
- 部署一個 Cloud Run 作業,為增量的價格更新執行輕量級的資料擴充。當白天有小型更新檔案送達時,透過監聽來自 Eventarc 的 Pub/Sub 訊息的 Cloud Workflows 來觸發此作業。
- 理由:使用 Workflows 的無伺服器容器 (Serverless containers) 為小型事件提供了低延遲、低維運的 orchestration,同時將繁重的轉換作業保留在批次處理中。
可靠性控制
- 為 Dataflow 和 Cloud Run 作業中的暫時性故障設定具備指數輪詢 (exponential backoff) 的重試機制,並將總重試時間限制在 DAG 的 SLA 內。在 Airflow 中設定每個任務的執行逾時 (execution timeouts) 與 on_failure 回呼;在 Workflows 中,設定 max_doublings 與 max_retry_duration。
- 理由:有界線的輪詢 (Bounded backoff) 能維持 SLA 並防止失控的重試。
安全性與最小權限原則
- 讓每個元件都在專用的服務帳戶 (service account) 下執行:Composer orchestrator SA、Dataflow worker SA、Cloud Run job SA。僅授予必要的角色:將擷取儲存貯體的 GCS 讀取權限授予 Dataflow、將目標資料集的 BigQuery dataEditor 角色授予相關服務,以及日誌的 Viewer 角色。將密鑰儲存在 Secret Manager 中,並在執行期參考它們。
- 理由:強制執行最小權限原則並隔離爆炸半徑 (blast radius)。
由元資料 (Metadata) 驅動的 Orchestration
- 維護一個 BigQuery 控制資料表,其中列出合作夥伴來源、檔案模式和目標資料集。在 DAG 執行期間,Airflow 會查詢此資料表,並使用動態任務對應 (dynamic task mapping) 來為每個合作夥伴生成對應的任務。
- 理由:新增合作夥伴變成了一項資料變更,而非程式碼變更,從而降低了部署風險。
可觀測性與警報
- 發送帶有 run_id 和 partner_id 的結構化日誌。為 DAG SLA 未達成、Dataflow 系統延遲和無效信件計數非空的情況建立警報政策。針對目標資料表的 BigQuery 插入作業,設定一個 Cloud Logging sink,並使用進階篩選器將該資料表的日誌導向一個 Pub/Sub 主題,供 Acme 的監控工具取用。
- 理由:細粒度的警報能夠在沒有雜訊的情況下實現快速分類處理 (triage)。
CI/CD 與環境推廣
- 使用 Terraform 管理基礎設施 (儲存貯體、Pub/Sub、Eventarc、Composer、Workflows、BigQuery 資料集)。使用 Cloud Build 來驗證 Airflow DAG 語法、執行單元測試,並部署到開發用的 Composer 環境。在 Dataform 斷言和整合測試通過後,使用參數化設定檔和手動批准關卡 (manual approval gates) 將其推廣到測試 (test) 和生產 (prod) 環境。
- 理由:宣告式、可重複的部署,以及跨環境的安全推廣。
Runbook 與復原
- 文件化重播特定日期的步驟:從物件版本管理中還原 CSV、針對該分割區重新執行 Dataflow 作業、MERGE 結果,並檢視 DLQ 記錄。包含一個「固定資料集重播 (fixed dataset replay)」程序,以便在出現差異時隔離轉換過程中的錯誤。
- 理由:冪等性設計和文件化的復原程序簡化了部分失敗的修復工作。
← 資料擷取、整合與遷移 · 所有領域 · 機器學習、AI 與資料服務 →
練習這些題目 → · 在 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.
通過考試 →