Lập lịch theo dữ liệu với Asset

Xây dựng Data Pipeline với Airflow

Volker Janz

Senior Developer Advocate at Astronomer

Vấn đề về thời điểm

 

  • Pipeline A nạp dữ liệu thô từ một API
  • Pipeline B dựng dashboard từ dữ liệu đó
  • Nếu cả hai chạy theo cron, B có thể chạy trước khi A xong
  • Thêm bộ đệm thời gian thì dễ vỡ

Vấn đề về thời điểm

Xây dựng Data Pipeline với Airflow

Asset là gì?

from airflow.sdk import Asset
sales_data = Asset("s3://bucket/sales/daily.parquet")
  • Một tham chiếu đến một phần dữ liệu, được định danh bằng tên duy nhất
  • Có thể gắn URI cho asset khi nó biểu diễn một thực thể dữ liệu cụ thể
  • Có thể là tệp, bảng cơ sở dữ liệu, hoặc bất kỳ dữ liệu nào
Xây dựng Data Pipeline với Airflow

Bên tạo và bên tiêu thụ

  • Task có thể cập nhật asset, tạo ra một asset event
  • DAG có thể được kích hoạt khi asset cập nhật bằng asset schedule

 

Luồng asset

Xây dựng Data Pipeline với Airflow

Bên tạo: phát tín hiệu với outlets

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

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

 

  • Khi thành công, Airflow ghi nhận asset đã được cập nhật
  • Bên tạo không cần biết về bất kỳ bên tiêu thụ nào
Xây dựng Data Pipeline với Airflow

Bên tiêu thụ: lập lịch theo asset

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

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

 

  • Không cần biểu thức cron, chỉ tham chiếu asset
  • Airflow tự kích hoạt bên tiêu thụ sau khi cập nhật
Xây dựng Data Pipeline với Airflow

Lập lịch có điều kiện

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

# Chờ cả HAI asset đều cập nhật @dag(schedule=(sales & inventory)) def full_report(): ...
# Kích hoạt khi BẤT KỲ asset nào cập nhật @dag(schedule=(sales | inventory)) def quick_refresh(): ...
Xây dựng Data Pipeline với Airflow

Xác minh asset với 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 hiển thị tất cả asset đã đăng ký
  • airflow assets details hiển thị metadata gồm cả updated_at
Xây dựng Data Pipeline với Airflow

Cùng luyện tập nào!

Xây dựng Data Pipeline với Airflow

Preparing Video For Download...