Google PDE: 訊息傳遞、事件擷取與即時服務 — 學習指南
屬於 Google Professional Data Engineer — 學習指南. 使用經過驗證的解答練習: Google 考試中心, 或參加限時模擬考試: ExamRoll.io.
總覽
Google Cloud 上的訊息傳遞、事件擷取與即時服務,其核心為用於解耦、持久傳輸的 Cloud Pub/Sub 和 Eventarc;用於具狀態串流處理的 Dataflow;以及如 BigQuery、Cloud Storage 和營運資料庫等匯集點 (sinks)。設計時若能考量到至少一次傳遞 (at-least-once delivery)、冪等性消費 (idempotent consumption) 及可觀測性 (observability),便能確保系統具備韌性,可在維持彈性擴展的同時,於故障、背壓 (backpressure) 和結構演變 (schema evolution) 的情況下仍能維持正確性。
使用 Pub/Sub 的核心訊息傳遞
- 主題 (Topics) 與訂閱 (subscriptions)
- 發布者 (Publishers) 將訊息傳送至一個主題;訂閱者 (subscribers) 透過訂閱來附加(多個訂閱者可以獨立地消費相同的訊息)。
- 訂閱類型:
- 拉取 (Pull):用戶端明確地拉取訊息;使用串流拉取 (streaming pull) 可獲得最高吞吐量並減少往返次數。
- 推送 (Push):Pub/Sub 透過 HTTPS 傳遞訊息;您的端點必須回傳 2xx 狀態碼以示確認。
- 匯出至 BigQuery:BigQuery 訂閱會將訊息傳遞至一個 BigQuery 資料表而無需撰寫程式碼;當酬載 (payloads) 符合宣告的結構 (schema) 且需要低延遲地將資料擷取至分析系統時,此為最佳選擇。
- 排序鍵 (Ordering keys)
- 在主題和訂閱上啟用訊息排序,以根據排序鍵接收依序傳遞的訊息。每個排序鍵的吞吐量是序列化的:每個鍵一次只有一則傳輸中的訊息,這可能會阻擋後續訊息;使用多個鍵(例如,hash(device_id))來擴展效能。
- 扇出 (Fan-out) 與重播 (replay)
- 為不同的消費者建立獨立的訂閱,以隔離工作負載和保留設定。
- 使用尋找 (Seek) 或快照 (snapshot) 功能,從特定時間戳或快照重播訊息,以進行復原和資料回填。
權衡取捨:
- 排序會降低每個鍵的平行處理能力和吞吐量;除非絕對必要,否則應停用排序。
- 推送簡化了用戶端程式碼,但引入了 HTTP 端點的擴展、安全性和退避機制 (backoff) 的考量;拉取則在高吞吐量下提供更多的控制權和穩定性。
傳遞語意、確認、保留與無法處理的信件
- 確認 (Acknowledgment) 與期限 (deadlines)
- 至少一次傳遞 (At-least-once delivery):可能會出現重複的訊息。
- 每次傳遞都有一個確認期限 (ack deadline)(預設為 10 秒)。在處理長時間執行的工作時,可延長此期限 (ModifyAckDeadline);未在期限前確認是導致重複推送傳遞最常見的原因。
- 否認 (Nack) 或期限到期會使訊息有資格被重新傳遞。
- 保留 (Retention)
- 未被確認的訊息會根據訂閱的確認期限被保留並重試;已被確認的訊息則可根據主題的訊息保留期限被保留,以供重播。設定保留時間時,應涵蓋您最長的服務中斷時間加上恢復所需的時間。
- 重試 (Retries)
- 拉取 (Pull):在確認期限到期後會重新傳遞;可使用流量控制限制來控制並行處理的數量。
- 推送 (Push):採用指數退避 (exponential backoff);只有 HTTP 2xx 狀態碼被視為成功。3xx/4xx/5xx 會觸發重試。實作冪等處理器 (idempotent handlers) 以容忍重複的請求。
- 無法處理的信件主題 (Dead-letter topics, DLTs)
- 為每個訂閱設定一個 DL 主題和最大傳遞嘗試次數,以隔離有問題的訊息 (poison messages)。
- 監控 DLQ 的數量;建立分流工作流程,並在修正後將訊息重新發布到主要主題。
範例:
undefined
傳遞語意摘要:
- Pub/Sub:至少一次傳遞,若啟用排序鍵,則在排序鍵內盡力排序 (best-effort ordering)。
- 匯集點 (Sinks):BigQuery 的 insert API 提供重複資料緩解機制 (insertId 或 Storage Write API 的串流偏移量),但消費者和寫入者仍應設計為冪等的。
結構、相容性與驗證
- Pub/Sub 結構 (Schemas)
- 原生支援 Avro 和 Protocol Buffers,並將結構集中儲存。
- 主題層級的結構設定:編碼 (Avro 或 Protobuf) 和強制執行 (無、僅驗證、或要求)。
- 生產者發布編碼後的酬載;當啟用強制執行時,Pub/Sub 會根據當前的結構進行驗證。
- 演變與相容性
- 使用向後相容的變更(新增可選欄位、在 Avro 中新增具預設值的欄位、在 Protobuf 中絕不重複使用標籤、避免移除或重新命名欄位)。
- 明確地對結構進行版本控制。對於破壞性變更,可同時發布到 v1 和 v2 主題,或新增一個版本欄位並據此進行路由。
- 生產者–消費者合約
- 消費者應忽略未知的欄位,並為遺失的欄位提供預設值。
- 在推廣至生產環境前,應在所有消費者之間測試結構的相容性;在預備 (staging) 環境的訂閱中,使用與生產環境相同的結構強制執行設定進行驗證。
簡短的 Avro 範例(摘錄):
undefined
事件驅動整合、Eventarc 與 Kafka 互通性
- Eventarc 與 CloudEvents
- Eventarc 使用 CloudEvents 規格,將來自 Google Cloud 服務 (以及透過 Pub/Sub 的自訂來源) 的事件路由到 Cloud Run、GKE 或 Workflows。像 type、source、subject 這類的屬性,能夠進行精細的篩選並提供可稽核性。
- 使用屬性篩選器來最小化扇出 (fan-out) 並減少下游的負載。
- 傳遞語意為「至少一次」(at-least-once);應盡可能讓處理程式 (handler) 具備冪等性 (idempotent) 且無狀態 (stateless)。
- Eventarc 觸發條件範例:
gcloud eventarc triggers create gcs-finalize-to-run
–destination-run-service=ingestor
–event-filters=“type=google.cloud.storage.object.v1.finalized”
–event-filters=“bucket=my-data-bucket”
–service-account=eventarc-sa@PROJECT_ID.iam.gserviceaccount.com - Kafka 互通性與託管式遷移
- Dataflow 範本可連接 Kafka <-> Pub/Sub,以進行分階段遷移。鏡像 (mirror) 主題並保留鍵值 (key);先切換消費者 (consumer),再切換生產者 (producer),或在轉換期間進行雙寫 (dual-write)。
- Pub/Sub Lite 提供分區 (partitioned)、容量預配 (capacity-provisioned) 的串流服務,具備基於鍵值的路由功能且成本較低;它是區域級/可用區級 (regional/zonal) 的服務,適合類似 Kafka 的工作負載,其中可預測的容量和每個分區的順序性是主要考量。
- 遷移考量事項:
- 順序性:將 Kafka 的鍵值 (key) 對應到 Pub/Sub 的排序鍵 (ordering key) 或 Lite 的分區 (partition)。
- 位移 (Offset):將位移作為訊息屬性攜帶以供診斷;遷移後,消費者不能再依賴 Kafka 的位移。
- 傳遞語意:接受「至少一次」(at-least-once);在下游強制執行冪等性。
- 結構描述 (Schema):將 Confluent Schema Registry 的定義遷移到 Pub/Sub 結構描述,或標準化採用 Protobuf/Avro,並遵循相容的演進規則。
串流擷取模式、吞吐量、擴展、安全性與維運
- 即時擷取模式
- Pub/Sub -> Dataflow -> BigQuery:使用 BigQuery Storage Write API sink 來達成高吞吐量與透過串流偏移量 (stream offsets) 實現的冪等性;將失敗的訊息路由到 dead-letter table 以供檢視。
- Pub/Sub -> Dataflow -> Cloud Storage:封存原始事件以供重新處理;使用視窗化、壓縮寫入來平衡成本與延遲。
- Pub/Sub -> 營運型資料庫:寫入 Bigtable 以進行低延遲查詢,寫入 Spanner 以進行強一致性交易,或根據工作負載需求寫入 Cloud SQL/Firestore。確保以唯一的事件 ID 為鍵值進行冪等的 upsert 操作。
- 至少一次 (At-least-once)、防止重複與冪等性
- 在每則訊息中攜帶唯一的 event_id 和 event_time;在生產者端強制使用 UUID。
- BigQuery 串流去重:設定 insertId 或使用帶有排序串流的 Storage Write API;查詢時仍需加上去重邏輯來防護。
- 查詢時去重範例:
undefined
- 對於 push endpoint,僅在成功處理後才回傳 2xx 狀態碼;否則預期訊息會被重新投遞。
- 訊息吞吐量、配額與擴展
- 發布者 (Publishers):批次處理訊息並重複使用連線;跨多個客戶端進行平行化。使用多個排序鍵 (ordering keys) 來擴展有排序需求的負載。
- 訂閱者 (Subscribers):偏好使用帶有流量控制 (flow control) 的 streaming pull (設定 max outstanding bytes/messages)。根據處理時間來設定 ack deadline 的長度,並在需要時延長。
- 隨著流量增長,監控並申請提高發布和訂閱吞吐量的配額;設計時應保留餘裕空間 (例如,預期峰值的 2 倍) 以吸收突發流量。
- 一致性與可用性
- BigQuery 串流對於查詢可見性是最終一致的;對於必須包含串流資料列的互動式查詢,應根據觀察到的延遲來等待 (例如,P50 可用性延遲的 2 倍),或在 Dataflow 中使用與浮水印 (watermark) 對齊的聚合,並查詢具體化後的結果。
- 安全性
- IAM:授予最小權限角色 (在 topic 上給予生產者 pubsub.publisher;在 subscription 上給予消費者 pubsub.subscriber)。為每個工作負載使用專屬的服務帳戶。
- Push 驗證:設定 push subscription,使其附加來自服務帳戶的 OIDC token;在 endpoint 上強制執行 audience 驗證。偏好使用 Cloud Run private endpoint,因其內建驗證與 TLS。
- 加密:Pub/Sub 會對傳輸中和靜態資料進行加密;在 topic 上使用 CMEK 以實現客戶管理的金鑰。應用 VPC Service Controls 來降低資料外洩的風險。若有需要,可對敏感的 payload 欄位使用客戶端加密。
- 延遲、重新投遞與訂閱者失敗的維運診斷
- 使用 Cloud Monitoring 進行監控:
- 使用 subscription/num_undelivered_messages 和 oldest_unacked_message_age 來監控積壓 (backlog)。
- 使用 expired_ack_deadline_count 來偵測因錯過 ack 而導致的重複訊息。
- 使用 publish_request_count 和 pull_request_count 來監控吞吐量。
- 透過將已知的資料集在 pipeline 中重放,並逐階段比較輸出,以找出有問題的轉換 (transform) 或 sink,來調查儀表板上遺失的事件。
- 對於 Dataflow 串流:
- 使用自動擴展並設定適當的 maxWorkers 來吸收來自多個來源的負載。
- 對於不相容的更新,應使用 drain 模式來部署 pipeline,以允許進行中的工作完成並防止資料遺失。
- 對於 BigQuery 的插入通知,可透過一個 sink 將篩選特定資料表的 Cloud Logging 稽核條目路由到一個 Pub/Sub topic 以進行告警。
- 使用 Cloud Monitoring 進行監控:
實務問題情境
Contoso Freight 公司需要一個全球性的即時事件平台,用來每分鐘從卡車擷取 10,000 筆 IoT 遙測訊息、豐富化事件、支援互動式分析,並在外部合作夥伴放置檔案時觸發工作流程。部分合作夥伴的 CSV 檔案包含格式錯誤的資料列,分析團隊必須能夠在不阻斷串流的情況下檢視這些錯誤。
- 建立核心訊息傳遞與 schema 層
- 行動:為遙測資料定義一個 Avro schema,並將其附加到一個名為
telemetry的 Pub/Sub topic,同時將 schema 強制執行設定為require。啟用訊息排序功能,並使用ordering_key = hash(device_id)來發布訊息。 - 理由:在 topic 層級強制執行 schema 能及早拒絕格式錯誤的事件。依裝置排序可在需要時支援有序處理,而雜湊 (hashing) 則能分散鍵值以維持吞吐量。
- 配置具備隔離與 dead-lettering 功能的 subscription
- 行動:為 Dataflow 建立一個名為
telemetry-stream-sub的 pull subscription,並設定一個 dead-letter topictelemetry-dlt及max_delivery_attempts=10。再新增一個 BigQuery subscriptiontelemetry-raw-bq,將原始事件存入一個時間分區的資料表,以供資料血緣 (lineage) 和重放之用。 - 理由:DLQ (Dead-Letter Queue) 能隔離有問題的訊息 (poison messages) 以供調查。一個獨立的 BigQuery subscription 提供了一個低維運成本的匯出路徑,用於保存原始事件,且與處理 pipeline 無關。
- 建立一個用於資料豐富化與寫入 sink 的 Dataflow 串流 pipeline
- 行動:使用帶有流量控制的 streaming pull 從
telemetry-stream-sub擷取資料。根據 schema 進行驗證,用參考資料進行豐富化,並計算視窗化聚合。使用 Storage Write API 將資料寫入 BigQuery,並指定一個具名串流 (named stream) 及insertId = event_id;每小時將原始備份寫入 Cloud Storage;將格式錯誤或失敗的記錄重新導向到一個 dead-letter BigQuery 資料表。 - 理由:Storage Write API 透過
insertId或串流偏移量 (stream offsets) 提供了高吞吐量、低延遲且具備冪等性的寫入。一個 dead-letter table 支援在不阻斷串流的情況下進行檢視,而 Cloud Storage 的封存則可用於重放。
- 處理分析中的重複資料與最終一致性
- 行動:對於必須排除重複資料的互動式查詢,在每筆記錄中發布
event_id和event_time,並使用一個去重複的 view:
undefined
根據觀察到的 BigQuery 串流可用性延遲 (例如,中位數延遲的兩倍) 在查詢前引入一個短暫的延遲。
- 理由:At-least-once 的傳遞模式需要冪等寫入和查詢時去重複。考慮到串流資料的可見性延遲,等待可以減少因資料尚在傳輸中而查詢不到的問題。
- 使用 Eventarc 整合合作夥伴的檔案放置
- 行動:設定 Eventarc,將
partner-dropsbucket 的 Cloud Storageobject.finalized事件路由到一個 Cloud Run 服務。該服務會啟動一個批次 Dataflow job 來將 CSV 載入 BigQuery,並將解析錯誤的資料傳送到一個 dead-letter table。 - 理由:Eventarc 提供了事件驅動的協作機制,可透過 CloudEvents 依 bucket 和物件前綴進行篩選。一個批次的 Dataflow job 能將格式錯誤的資料列分離出來進行分析,同時迅速載入正確的資料。
- 保護平台安全
- 行動:使用不同的服務帳戶:生產者在
telemetrytopic 上獲得pubsub.publisher權限;Dataflow worker SA 在telemetry-stream-sub上獲得pubsub.subscriber權限,以及對目標 BigQuery dataset 和 Cloud Storage 的寫入權限;Eventarc trigger 使用一個專屬的 SA,並具備對 Cloud Run 的 invoker 權限。在telemetrytopic 和 BigQuery dataset 上啟用 CMEK。若有使用 push endpoint,則設定 OIDC 和 audience 檢查。 - 理由:最小權限 IAM 和 CMEK 符合安全性與合規性要求;經過驗證的傳遞可防止偽造攻擊。
- 可靠地維運與擴展
- 行動:設定 Dataflow 自動擴展,並給予一個充裕的
maxWorkers值以吸收峰值流量。監控subscription/oldest_unacked_message_age和expired_ack_deadline_count;當超過閾值時發出告警。對於會破壞相容性的 pipeline 變更,應使用drain模式部署以避免訊息遺失。如果延遲增加,應提高訂閱者的平行處理能力,並根據處理時間按比例延長 ack deadline。 - 理由:主動監控能及早偵測到延遲和重新投遞。自動擴展和調整過的 ack deadline 能防止重複訊息風暴。在升級期間使用
drain模式能保留進行中的訊息。
此設計提供了一個具備彈性、安全且可觀測的即時擷取方案,並整合了事件驅動的批次處理,支援重複容錯和 schema 演進,同時在隔離錯誤資料以進行針對性修復的基礎上,提供快速的分析能力。
← 使用 Dataflow 與 Apache Beam 進行串流處理 · 所有領域 · 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.
通過考試 →