Asset Partitions를 활용한 파티션 인식 스케줄링

Airflow로 데이터 파이프라인 구축하기

Volker Janz

Senior Developer Advocate at Astronomer

ds만으로는 부족할 때

 

  • {{ ds }}는 각 실행을 날짜에 연결하므로 시간 스케줄 기반 DAG에 적합합니다
  • asset 트리거 DAG에서는 dsNone으로 설정되어 사용할 수 없습니다
  • 다운스트림 DAG는 어떤 날짜가 방금 업데이트되었는지 알 수 없습니다
  • Asset Partitions는 업스트림에서 다운스트림으로 partition_key를 전파합니다

업스트림 스케줄 실행이 ds가 None인 다운스트림 asset 스케줄 실행을 트리거하는 모습

Airflow로 데이터 파이프라인 구축하기

CronPartitionTimetable

from airflow.sdk import dag, CronPartitionTimetable

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

 

  • 스케줄 실행partition_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로 데이터 파이프라인 구축하기

파티션 asset 이벤트

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 = 파티션 asset 이벤트
  • 이벤트는 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}")

 

  • 파티션 asset 이벤트에서 트리거됩니다
  • 다운스트림 실행은 파티션 키를 상속합니다
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

$$

  • 시간 기반 키: StartOfHourMapper, StartOfWeekMapper, StartOfMonthMapper
  • 비시간 기반 키: 지역, 부서 등에는 AllowedKeyMapper 사용
Airflow로 데이터 파이프라인 구축하기

전체 흐름

전체 파티션 흐름

  • 업스트림: CronPartitionTimetable이 각 실행에 파티션 키를 부여합니다
  • Outlet 태스크가 해당 키를 포함한 파티션 asset 이벤트를 발행합니다
  • StartOfDayMapper가 타임스탬프를 날짜 형식으로 정규화합니다
  • 다운스트림: PartitionedAssetTimetable이 트리거되어 매핑된 키를 상속합니다
Airflow로 데이터 파이프라인 구축하기

연습해 봅시다!

Airflow로 데이터 파이프라인 구축하기

Preparing Video For Download...