使用 Airflow 构建数据流水线
Volker Janz
Senior Developer Advocate at Astronomer

from airflow.sdk import Asset
sales_data = Asset("s3://bucket/sales/daily.parquet")

sales_data = Asset("s3://bucket/sales/daily.parquet")@task(outlets=[sales_data]) def write_sales(): # Write data to S3 ...
sales_data = Asset("s3://bucket/sales/daily.parquet")@dag(schedule=[sales_data]) def build_dashboard(): ...
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 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 构建数据流水线