Airflow로 데이터 파이프라인 구축하기
Volker Janz
Senior Developer Advocate at Astronomer
{{ ds }}는 각 실행을 날짜에 연결하므로 시간 스케줄 기반 DAG에 적합합니다ds가 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로 데이터 파이프라인 구축하기