使用 Assets 进行数据感知调度

使用 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")
  • 对数据的一个引用,用唯一名称标识
  • 当表示具体数据实体时,可为资产附加一个 URI
  • 可以是文件、数据库表,或任意数据
使用 Airflow 构建数据流水线

生产者与消费者

  • 任务可更新资产,从而产生资产事件
  • Dags 可用资产日程在资产更新时触发

 

资产流转

使用 Airflow 构建数据流水线

生产者:用 outlets 发出信号

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

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

 

  • 成功后,Airflow 会记录该资产已更新
  • 生产者无需了解任何消费者
使用 Airflow 构建数据流水线

消费者:基于资产调度

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

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

 

  • 无需 cron 表达式,只用资产引用
  • 资产更新后,Airflow 会自动触发消费者
使用 Airflow 构建数据流水线

条件调度

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

# 同时等待两个资产更新 @dag(schedule=(sales & inventory)) def full_report(): ...
# 任一资产更新即触发 @dag(schedule=(sales | inventory)) def quick_refresh(): ...
使用 Airflow 构建数据流水线

用 CLI 验证资产

$ 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 显示所有已注册的资产
  • airflow assets details 显示包含 updated_at 在内的元数据
使用 Airflow 构建数据流水线

让我们一起练习吧!

使用 Airflow 构建数据流水线

Preparing Video For Download...