以 Asset 進行資料感知排程

使用 Airflow 建置資料管線

Volker Janz

Senior Developer Advocate at Astronomer

時序問題

 

  • Pipeline A 從 API 載入原始資料
  • Pipeline B 用該資料建立儀表板
  • 若兩者都用 cron,B 可能在 A 完成「之前」執行
  • 加上時間緩衝很脆弱

時序問題

使用 Airflow 建置資料管線

什麼是 Asset?

from airflow.sdk import Asset
sales_data = Asset("s3://bucket/sales/daily.parquet")
  • 對某筆資料的參照,以唯一名稱識別
  • 當代表具體資料實體時,可為 asset 附上 URI
  • 可以是檔案、資料庫資料表,或任何資料
使用 Airflow 建置資料管線

生產者與消費者

  • 工作可更新 asset,並產生asset 事件
  • Dag 可用asset 行程在 asset 更新時觸發

 

Asset 流程

使用 Airflow 建置資料管線

生產者:用 outlets 發出訊號

sales_data = Asset("s3://bucket/sales/daily.parquet")

@task(outlets=[sales_data]) def write_sales(): # Write data to S3 ...

 

  • 成功後,Airflow 會記錄該 asset 已更新
  • 生產者無需知道任何消費者
使用 Airflow 建置資料管線

消費者:依 asset 排程

sales_data = Asset("s3://bucket/sales/daily.parquet")

@dag(schedule=[sales_data]) def build_dashboard(): ...

 

  • 不需 cron 表達式,只要asset 參照
  • 更新後,Airflow 會自動觸發消費者
使用 Airflow 建置資料管線

條件式排程

sales = Asset("s3://bucket/sales.parquet")
inventory = Asset("s3://bucket/inventory.parquet")

# 同時等待兩個 asset 更新 @dag(schedule=(sales & inventory)) def full_report(): ...
# 任一 asset 更新即觸發 @dag(schedule=(sales | inventory)) def quick_refresh(): ...
使用 Airflow 建置資料管線

用 CLI 驗證 assets

$ airflow assets list
name                                  | group    | uri
s3://data-lake/sales/daily.csv        | asset    | s3://data-lake/sales/daily.csv
$ airflow assets details --name "s3://data-lake/sales/daily.csv"
property_name    | property_value
=================+=====================                                         
name             | s3://data-lake/sales/daily.csv
                ...
updated_at       | 2026-04-14T10:00:50.820041Z
  • airflow assets list 會列出所有已註冊的 asset
  • airflow assets details 會顯示包含 updated_at 的中繼資料
使用 Airflow 建置資料管線

一起來練習吧!

使用 Airflow 建置資料管線

Preparing Video For Download...