Day 14 我處理的是設備事件進入 Pub/Sub 後,如何確認、重試與隔離。接下來,我想知道訊息持續寫入 BigQuery 時,怎樣才能看出設備讀數隨時間變化。Day 13 已經提醒我,設備識別資訊、讀數單位和重複事件都要先有規則;到了串流資料,還多了一個問題:一筆讀數應算在設備測量的時刻,還是系統收到它的時刻?
假設設備暫時離線,恢復連線後才送出先前測得的溫度。資料表現在才新增一列,不表示溫度是現在測得的。我若直接拿資料到達的順序解讀設備狀態,可能把較早發生的變化誤認成剛發生。Day 13 檢查每日統計的品質;這裡進一步追問事件抵達後,時間歸屬該怎麼定義。
我先建立目標表格,確認要寫入的 JSON 與欄位結構相容,再設定 Pub/Sub 服務代理的寫入權限、建立 BigQuery 訂閱,最後發布訊息並查表驗證。BigQuery 訂閱屬於匯出型訂閱,能把訊息寫入既有表格,不必另外執行接收程式。它幫我完成的是持續落地;查到資料後,仍要檢查欄位是否按預期保存。
在設備事件的設計中,我會先決定設備識別資訊、測量時間與讀數放在哪些欄位。BigQuery 訂閱可以依設定使用主題或表格的 Schema 對應資料,也能將訊息本文寫入 data 欄位;選項不同,目標表格需要的欄位也不同。不能因為訊息本文是 JSON,就以為它會自動拆成我想要的分析欄位。
訊息落地後,我用 SQL 依事件時間將資料歸入固定區間,檢查各區間的彙總。撰寫查詢時,我也讓 AI 協助產生語句;但核對後發現它設定的區間與我的要求不同。查詢能執行,只能證明語法過關,不能證明時間邊界、分組欄位和業務定義都正確。
BigQuery 的 TIMESTAMP_BUCKET 可以取得某個時間值所屬時間桶的起點。我會先確認測量時間已轉成正確的 TIMESTAMP,再選定區間寬度與起點,依設備分組查詢。這是在當下對已寫入表格的資料重新計算;下次有新事件或補送事件進來,必須再查一次,才會得到包含新資料的結果。它並沒有替我建立持續執行、定時輸出視窗結果的工作。
設備記錄的測量時間回答「現場何時發生」;Pub/Sub 的發布時間回答「訊息何時送進系統」;若另外記錄入表時間,才能回答「資料何時寫入目標表」。設備離線補送時,這幾個時間可能相距很遠。若要用發布時間分析,我也必須確認訂閱已設定寫入相關訊息中繼資料;若想比較入表延遲,就得先設計如何取得入表時間,不能假設每張表都有現成欄位。
以設備狀態為目的,我通常會先考慮測量時間,同時保留足夠的接收線索,才能查出資料為何晚到。Day 13 的品質檢查在這裡仍有用:時間能否解析、時區是否一致、設備識別是否完整、同一事件是否重複計入,都會直接改變分組結果。更短的區間無法補救錯誤的時間或重複資料。
假設我上午查過某段時間的設備讀數,下午才收到設備補送的事件。再執行相同 SQL,先前區間的筆數或平均值可能改變。對查詢而言,這是表格新增資料後的合理結果;但對監控畫面而言,我得先決定它顯示的是暫定值,還是某個時點後不再修正的結果。
Dataflow 的串流文件用 watermark 表示系統對事件時間進度的估計;當事件在 watermark 越過視窗終點後才到達,它就屬於遲到資料。trigger(觸發規則)則決定何時輸出結果。這些概念讓我能討論「要等多久、要不要提前顯示、晚到時如何修正」,但前面的 BigQuery SQL 沒有實作 watermark 或 trigger,不能把回溯查詢說成持續處理的串流視窗。
表格暫時沒有新讀數,不一定是設備停止測量,也可能是訂閱寫入失敗。Google Cloud 的故障排除文件列出表格不存在、寫入權限不足與 Schema 不相容等情況。我會檢查訂閱狀態和錯誤原因,再核對目標表的資料;對無法寫入的訊息,也可以規劃死信路徑,留下後續處理線索。
落地層與品質層仍要分開驗收。前者確認訊息是否持續寫入,後者確認寫入後的資料符合 Day 13 所訂的規則。若同一筆測量被重新發布,BigQuery 訂閱的至少一次投遞語意也不能替我完成業務上的去重;只檢查表格有沒有新列,會漏掉這類問題。
做完這段操作,我會先問資料使用者:可以在需要時重查,並接受舊區間的結果因補送而改變嗎?還是監控系統必須持續取得結果,並明確知道何時發布、何時修正?如果是前者,BigQuery 訂閱加上經驗證的 SQL 已能提供一條清楚的路徑。如果是後者,我就得進一步設計視窗輸出和遲到資料的處理方式。
Day 13 讓我學會檢查結果是否符合資料定義;這次則讓我看見,結果何時算得出來、何時可交付,也是定義的一部分。下一步才是判斷:這個需求是否值得增加一條持續運行的 Dataflow 管道。