iT邦幫忙

2026 iThome 鐵人賽

DAY 16
0
AI 自動化

從 Cloud 到 AI:GCP 雲端資料與 AI 實戰之路系列 第 16 篇

Day 16|什麼時候才需要 Dataflow?

  • 分享至 

  • xImage
  •  

1. BigQuery 已能收資料,為什麼還要增加一條管道?

Day 15 我讓 Pub/Sub 訊息寫入 BigQuery,再對已落地資料依事件時間查詢。那條路徑能保存事件,也能讓我重新計算補送資料所影響的區間。但如果結果需要持續輸出,或訊息在入表前必須經過較複雜的處理,就得重新檢查處理流程的能力。

我原本容易把 Dataflow 想成「串流架構的下一站」,彷彿資料進了 Pub/Sub 就應接上一個工作。實際增加一條管道後,我更在意的是:它承擔了哪項原本做不到的工作?我又要為啟動、監控與停止負起哪些責任?

2. 我先用範本建立第二條寫入路徑

操作時,我準備另一個訂閱與 BigQuery 目標表,設定工作區域、暫存資源和權限,再啟動 Pub/Sub to BigQuery 範本,觀察工作狀態與寫入結果。官方文件說明,這個範本讀取 Pub/Sub 的 JSON 訊息並寫入 BigQuery,也可選用 JavaScript 使用者定義函式(UDF)處理輸入;目標表須事先存在,資料須符合其 Schema。

我在這一步驗證的是一條持續搬運訊息的路徑。假設設備先後回報 24°C 和 26°C,範本會將兩筆讀數寫入 BigQuery,但不會直接輸出這段時間的平均溫度。Day 15 的 SQL 可以在我執行查詢時算出 25°C;如果設備稍後補送一筆測量時間屬於同一區間的 28°C,再次查詢時,平均值就會變成 26°C。

若希望資料持續進來時,系統也持續輸出各時間區間的統計,就需要另外設計處理管道:依設備測量時間分組、決定何時輸出結果,並規定補送資料是否修正先前的數字。Dataflow 可以執行這樣的管道。

3. 工作啟動失敗,先讀錯誤再改設定

我曾在啟動工作時遇到資源供應問題。當時先查看工作訊息與區域設定,調整後使用 Streaming Engine 完成了那次啟動。這段經驗讓我知道「參數填完」與「工作真的在執行」是兩個驗收點;一次成功的排錯方法,也需要放回當時的區域及資源條件理解。

Google Cloud 說明,區域或可用區的資源不足可能使 Dataflow 工作無法啟動;指定區域並使用 Streaming Engine,有助於利用區域部署能力。下次遇到啟動失敗,我會先核對工作記錄、區域、執行身分及所需資源,再決定調整哪項設定。這樣才能辨別是資源供應、權限還是輸入輸出設定造成的問題。

4. 兩條寫入路徑,要各自驗收

原有的 BigQuery 訂閱與新建的 Dataflow 訂閱,會各自接收其建立後符合條件的訊息,並寫入各自指定的表。比較結果時,我得先確認兩個訂閱何時開始接收、讀取的是哪些發布事件,以及工作是否有成功寫入。只比兩張表的列數,不能判斷哪一條路徑比較正確。

對設備事件,我還會分開看三件事:一則 Pub/Sub 訊息因處理失敗而重新投遞;來源將同一次測量重新發布成另一則訊息;兩個不同訂閱各自取得發布的訊息。它們都可能讓我在不同地方看見相似內容,卻有不同原因。Day 13 訂下的事件識別與去重規則,因此仍是分析結果可信的必要條件。BigQuery 訂閱本身採至少一次投遞,也不能只憑訂閱或表格筆數宣稱業務事件已達成精確一次。

5. 我用處理需求,而不是產品名稱作選擇

如果設備事件只需寫入表格,再以 SQL 查詢和驗證,BigQuery 訂閱已提供直接路徑。它也能附加輕量的單則訊息轉換,並為無法寫入的訊息設定死信主題。因此「需要簡單改寫欄位」或「需要留下寫入失敗的訊息」,本身還不足以說明必須啟動 Dataflow。

我會先問:處理一筆讀數時,需不需要等待其他讀數?例如把每筆溫度從攝氏換算成華氏,只要看當下這筆訊息;要持續計算某台設備在一段時間內的平均溫度,就得收集多筆讀數,並決定何時輸出、晚到的讀數要不要重新計算。前一種是單則訊息轉換;後一種是跨訊息的彙總。

因此,當需求涉及跨訊息計算、依事件時間持續產出結果,或較複雜的入表前轉換,我才會評估設計 Dataflow 管道。但選擇 Dataflow 後,仍要確認實際啟動的是哪一條管道:Pub/Sub to BigQuery 範本能完成文件列出的搬運與可選轉換;要持續計算平均溫度,還需要另外設計分組、輸出與補送資料的處理規則。

6. 串流工作要有人監控,也要有人結束

Dataflow 工作持續運行時,我要確認它是否仍在讀取、是否積壓、目標表是否持續寫入。工作區域、節點與計費模式也會影響成本,不能根據一次短暫測試推定正式運行的費用。Google Cloud 的成本文件建議依工作型態檢視資源配置;我會連同 Pub/Sub 與 BigQuery 的使用情況一起評估整條路徑。

結束工作時,cancel 會立即停止處理,可能影響已讀取但尚未寫完的資料;drain 則停止接收新資料,嘗試完成在途處理。我會依目的選擇,並確認工作最後狀態及目標表的寫入情形。若將來設計了視窗管道,停止時還得另外檢查尚未到期的區間會如何輸出;目前的搬運範本不能拿來證明這種行為。工作結束後,訂閱、表格和暫存資源也要逐一確認。

7. 多一條管道,必須多一個明確理由

走完這段操作,我現在會先說明需求,再選處理方式:訊息入口由 Pub/Sub 管理;單純落地及回溯查詢,可先用 BigQuery 訂閱與受驗證的 SQL;需要持續計算或較複雜的轉換時,再設計 Dataflow 管道,並把監控、成本與停止方式算進方案。

Day 13 到 Day 15 反覆追問資料從哪裡來、何時發生、是否通過檢查。Day 16 讓我再加上一問:每增加一條持續運行的路徑,我能否說清它提供了什麼結果,以及出錯時要由哪個環節負責?

參考資料

  1. Google Cloud, BigQuery subscriptions
  2. Google Cloud, Pub/Sub to BigQuery template
  3. Google Cloud, Dataflow regions
  4. Google Cloud, Stop a running Dataflow pipeline
  5. Google Cloud, Best practices for Dataflow cost optimization

上一篇
Day 15|資料進了 BigQuery,何時才能算進統計?
下一篇
Day 17|設備讀數之外,我如何理解操作人員的文字回報?
系列文
從 Cloud 到 AI:GCP 雲端資料與 AI 實戰之路 共 17 篇
圖片
  熱門推薦
圖片
{{ item.channelVendor }} | {{ item.webinarstarted }} |
{{ formatDate(item.duration) }}
直播中

尚未有邦友留言

立即登入留言