你的團隊有一條管線會把每日的銷售資料寫入 s3://data-lake/sales/daily.csv,而且有多條下游管線仰賴它。與其讓它們用計時器輪詢或猜測檔案何時就緒,你將加上一個 Asset outlet,讓這個任務在完成時立刻發出資料已更新的訊號。此訊號會讓其他 Dags 能夠直接依據該筆資料自行排程。
s3://data-lake/sales/daily.csv
完成程式碼後,先執行該檔案,接著使用 Airflow CLI 驗證該資產已被註冊並且已更新。
本練習屬於課程
將理論付諸實踐,立即體驗我們的互動練習
你將先認識 Airflow 的元件,使用 TaskFlow API 撰寫你的第一個 Dags,並透過 XCom 在任務之間傳遞資料。
接著,你會用動態任務對映平行執行任務,使用 Assets 以資料為基準排程 Dags,並加入人工核准步驟。
當前練習
本章將透過重試與回呼處理失敗情況,使用可延後的感測器節省資源,並以三個層級測試你的 Dags。
在最後一章,你會在 DuckDB 上建置 SQL ETL 管線,使用 Asset Partitions 加入分割區感知的排程,並內嵌資料品質檢查。