使用 Asset Partitions 的分割感知排程

使用 Airflow 建置資料管線

Volker Janz

Senior Developer Advocate at Astronomer

當 ds 不夠用時

 

  • {{ ds }} 會把每次執行綁到某個日期,適用於時間排程的 Dag
  • 但由資產觸發的 Dag 其 ds 會是 None,等同不可用
  • 下游 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...