Lập lịch theo phân vùng với Asset Partitions

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

Volker Janz

Senior Developer Advocate at Astronomer

Khi ds là không đủ

 

  • {{ ds }} gắn mỗi lần chạy với một ngày, phù hợp cho Dags lập lịch theo thời gian
  • Nhưng với Dags kích hoạt bởi asset, dsNone, tức là không có sẵn
  • Dag downstream không biết ngày nào vừa được cập nhật
  • Asset Partitions truyền partition_key từ upstream xuống downstream

Chạy upstream theo lịch kích hoạt downstream theo asset, nơi ds là None

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

CronPartitionTimetable

from airflow.sdk import dag, CronPartitionTimetable

@dag(schedule=CronPartitionTimetable("0 0 * * *", timezone="UTC"))
def sales_pipeline():
    ...

 

  • Mỗi lần chạy được lập lịch sẽ tự động có partition_key
  • Chạy thủ công sẽ không phân vùng trừ khi cung cấp key
Xây dựng Data Pipeline với Airflow

Truy cập partition key

Trong tác vụ Python:

@task
def process_data(dag_run=None):
    partition_key = dag_run.partition_key
    print(f"Processing: {partition_key}")

Trong mẫu SQL:

DELETE FROM daily_summary
WHERE order_date = '{{ dag_run.partition_key[:10] }}';

 

  • Có sẵn trong mọi tác vụ của một lần chạy được phân vùng
Xây dựng Data Pipeline với Airflow

Sự kiện asset có phân vùng

from airflow.sdk import dag, task, Asset, CronPartitionTimetable

daily_sales = Asset("daily_sales")

@dag(schedule=CronPartitionTimetable("0 0 * * *", timezone="UTC"))
def sales_pipeline():
    @task(outlets=[daily_sales])
    def load_data(**context):
        ...

 

  • Kết hợp outlet của tác vụ + CronPartitionTimetable = sự kiện asset có phân vùng
  • Sự kiện mang cùng partition key như lần chạy Dag
  • Báo hiệu phân vùng nào đã sẵn sàng, không chỉ là dữ liệu đã cập nhật
Xây dựng Data Pipeline với Airflow

PartitionedAssetTimetable

from airflow.sdk import dag, task, Asset, PartitionedAssetTimetable

daily_sales = Asset("daily_sales")

@dag(schedule=PartitionedAssetTimetable(assets=daily_sales))
def sales_report():
    @task
    def generate_report(dag_run=None):
        partition_key = dag_run.partition_key
        print(f"Report for: {partition_key}")

 

  • Chỉ kích hoạt trên các sự kiện asset có phân vùng
  • Lần chạy downstream kế thừa partition key
Xây dựng Data Pipeline với Airflow

StartOfDayMapper

from airflow.sdk import (
    dag, PartitionedAssetTimetable,
    StartOfDayMapper,
)

@dag(schedule=PartitionedAssetTimetable(
    assets=daily_sales,
    partition_mapper_config={
        daily_sales: StartOfDayMapper()
    },
))
def sales_report():
    ...

 

Trước khi mapping: 2026-04-23T00:00:00

Sau StartOfDayMapper: 2026-04-23

$$

  • Khóa theo thời gian: StartOfHourMapper, StartOfWeekMapper, StartOfMonthMapper
  • Khóa phi thời gian: AllowedKeyMapper cho vùng, phòng ban
Xây dựng Data Pipeline với Airflow

Luồng đầy đủ

Luồng phân vùng đầy đủ

  • Upstream: CronPartitionTimetable gắn partition key vào mỗi lần chạy
  • Tác vụ outlet phát ra sự kiện asset có phân vùng với key đó
  • StartOfDayMapper chuẩn hóa timestamp thành ngày
  • Downstream: PartitionedAssetTimetable kích hoạt và kế thừa key đã map
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...