使用 Asset Partitions 的分区感知调度

使用 Airflow 构建数据流水线

Volker Janz

Senior Developer Advocate at Astronomer

当 ds 不够用时

 

  • {{ ds }} 将每次运行绑定到日期,适用于按时间调度的 Dag
  • 资产触发的 Dag 中 dsNone,不可用
  • 下游 Dag 无法得知更新的是哪天
  • Asset Partitions 会将上游的 partition_key 传递到下游

上游定时运行触发下游按资产调度的运行,此时 ds 为 None

使用 Airflow 构建数据流水线

CronPartitionTimetable

from airflow.sdk import dag, CronPartitionTimetable

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

 

  • 每次定时运行都会自动获得一个 partition_key
  • 手动运行不会分区,除非显式提供 key
使用 Airflow 构建数据流水线

访问分区键

在 Python 任务中:

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

在 SQL 模板中:

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

 

  • 在分区运行中的每个任务都可用
使用 Airflow 构建数据流水线

分区资产事件

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):
        ...

 

  • 任务 outlet + CronPartitionTimetable = 分区资产事件
  • 事件携带与 Dag 运行相同的分区键
  • 标明哪个分区已就绪,而不只是数据已更新
使用 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}")

 

  • 仅在分区资产事件上触发
  • 下游运行会继承该分区键
使用 Airflow 构建数据流水线

StartOfDayMapper

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

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

 

映射前: 2026-04-23T00:00:00

使用 StartOfDayMapper 后: 2026-04-23

$$

  • 时间型键StartOfHourMapperStartOfWeekMapperStartOfMonthMapper
  • 非时间型键:地区、部门可用 AllowedKeyMapper
使用 Airflow 构建数据流水线

完整流程

完整分区流程

  • 上游:CronPartitionTimetable 为每次运行附加分区键
  • Outlet 任务发出带该键的分区资产事件
  • StartOfDayMapper 将时间戳规范为日期
  • 下游:PartitionedAssetTimetable 被触发并继承映射后的键
使用 Airflow 构建数据流水线

让我们一起练习吧!

使用 Airflow 构建数据流水线

Preparing Video For Download...