使用 Airflow 建置資料管線
Volker Janz
Senior Developer Advocate at Astronomer
{{ ds }} 會把每次執行綁到某個日期,適用於時間排程的 Dagds 會是 None,等同不可用partition_key 傳遞到下游
from airflow.sdk import dag, CronPartitionTimetable
@dag(schedule=CronPartitionTimetable("0 0 * * *", timezone="UTC"))
def sales_pipeline():
...
partition_key在 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] }}';
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):
...
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}")
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
$$
StartOfHourMapper、StartOfWeekMapper、StartOfMonthMapperAllowedKeyMapper 可用於地區、部門
CronPartitionTimetable 為每次執行附加分割鍵StartOfDayMapper 將時間戳正規化為日期PartitionedAssetTimetable 被觸發並繼承對映後的鍵使用 Airflow 建置資料管線