題幹與適用場景
上游工作每天產生多個資料資產,下游報表和品質檢查希望在依賴資產更新後自動執行。請說明如何用 Airflow 的 asset-aware scheduling 建模生產者、消費者、分割區和失敗復原,並解釋它與時間排程、外部感測器的邊界。
面試官考察點
- 是否把 asset 視為可由工作更新的邏輯資料依賴,而不是簡單的檔名或 cron 標籤。
- 是否能區分 asset 事件、DAG 時間表、運算式組合和事件觸發器。
- 是否考慮重複事件、延遲資料、分割區粒度、資料品質閘門和回填影響。
- 是否能說明監控、權限、冪等、失敗重試與暫停復原策略。
回答前需要釐清的問題
- 資產代表整張表、某個分割區還是物件儲存路徑?更新事件是否帶有分割區和批次資訊?
- 下游需要所有上游資產都更新,還是任一資產更新即可?是否存在跨 DAG 的資產運算式?
- 事件到達但資料品質檢查失敗時,應阻止消費 DAG、重試生產者,還是允許人工放行?
- 需要補歷史資料、重播事件或相容現有 cron 排程嗎?重複觸發的代價是什麼?
30 秒回答框架
我會先把業務資料產品建模成資產及其更新契約,再讓生產 DAG 在成功寫入並完成品質檢查後發出資產更新事件。消費 DAG 使用資產依賴或運算式決定觸發條件,而不是用更短的 cron 猜測資料是否就緒。接著定義分割區、冪等、重複事件、延遲更新和回填規則,並監控事件延遲、待觸發 DAG、失敗重試和資料新鮮度。若事件來自外部系統,則評估 Airflow event-driven trigger 的可復原性與安全邊界。
分步驟深入解答
1. 建立資產契約
在 Airflow Assets 官方文件中,資產是 DAG 之間共享的資料依賴。生產工作應在資料成功提交且契約滿足後更新資產;僅建立暫存檔、工作開始或部分寫入不應被當成可消費事件。資產 URI、所有權和分割區粒度需要穩定,否則下游無法判斷新舊資料。
2. 選擇觸發邏輯
消費 DAG 可以依賴一個或多個資產。根據 Asset-Aware Scheduling 官方文件,多個資產的邏輯組合要表達「全部更新」或「任一更新」等業務條件;時間表仍可作為獨立限制。設計時應明確事件和時間條件同時滿足時的行為,避免把運算式誤當成資料過濾器。
3. 處理分割區與重複
資產事件不會自動取代資料分割區水位。事件應攜帶可追蹤的批次或分割區資訊,消費工作用水位表、唯一鍵或交易寫入保證冪等。重複事件、工作重試和排程器復原都可能再次評估依賴;下游應能安全重跑,而不是依賴「每個事件只來一次」的假設。
4. 外部事件、品質與復原
如果更新來自訊息佇列或外部系統,可參考 Event-driven scheduling 官方文件選擇事件觸發器,但要核對觸發器是否支援可復原檢查、連線憑據和資源釋放。資料品質失敗時保留資產未就緒狀態,告警並重試或人工放行;回填時明確是否發出新事件、是否隔離歷史分割區,以及如何防止消費 DAG 與即時更新互相覆寫。
高品質示範回答
我會先定義資產契約:每個資產的 URI、分割區鍵、生產者、品質門檻和更新批次。生產 DAG 只有在原子寫入和品質檢查通過後才更新資產,消費 DAG 用資產依賴或運算式表達「全部上游完成」與「任一上游完成」,必要時再疊加時間表。事件中記錄分割區和批次,消費端用水位和冪等鍵處理重複、重試與排程器復原。外部更新採用受控的事件觸發器並限制憑據和資源生命週期。監控事件延遲、待觸發工作、資料新鮮度和失敗重播;回填與即時路徑使用不同批次邊界,防止歷史重播覆蓋最新結果。
常見錯誤
- 把 asset 當作任意檔案路徑,未定義誰在什麼時刻擁有更新權限。
- 只縮短 cron 間隔,卻沒有資料就緒、分割區和品質門檻。
- 認為資產事件天然只觸發一次,忽略重試、重複發布和排程器復原。
- 把資產依賴運算式當成資料內容過濾器,遺漏實際的分割區水位判斷。
- 外部事件觸發器長期持有連線或憑據,卻沒有逾時、釋放和告警策略。
- 回填時直接複用即時 DAG,導致歷史事件重複覆寫線上結果。
追問及應對
如果兩個上游資產到達時間不同怎麼辦?
明確下游是等待全部資產,還是允許先處理部分結果。等待全部資產時記錄每個分割區的到達水位和逾時告警;允許部分結果時輸出版本號,讓下游知道結果仍可被後續資產更新。
如何測試重複事件?
在測試環境重複發布同一資產更新、重試生產工作並重啟排程器,驗證消費工作的唯一鍵、水位和交易邊界。檢查結果筆數、版本號和外部副作用,確保重跑不會重複發送通知或扣費。
什麼時候保留外部感測器?
當依賴系統無法發布 Airflow 資產事件、需要輪詢受控介面,或遷移期間必須相容舊契約時,可以暫時使用感測器。應記錄遷移期限與輪詢成本,最終讓上游提供可驗證的資產更新訊號。