← 학습 카테고리

Learn

Airflow

151개 모듈 · 현재 41번째

Airflow 모듈 41/151 airflow-learn-41

Asset 파티션 기초 (partition_key와 PartitionedAssetTimetable)

Apache Airflow Official Documentation (in-repo snapshot) — Apache Software Foundation authoring-and-scheduling/assets.rst - Asset partitions (Pre-determined vs runtime partitioning ~ 기본 매퍼 부분, versionadded 3.2.0)

이 모듈을 다 읽으면

  • partition_key가 무엇이고 왜 파티션 인지 Dag만 트리거하는지 설명할 수 있다
  • 사전 결정(pre-determined) 파티셔닝과 런타임(runtime) 파티셔닝을 언제 각각 선택해야 하는지 판단할 수 있다
  • CronPartitionTimetable과 PartitionedAssetTimetable의 역할 분담을 설명할 수 있다
  • IdentityMapper·StartOf*Mapper·ProductMapper·AllowedKeyMapper·FixedKeyMapper의 동작 차이를 구분할 수 있다

Asset 이벤트에 partition_key를 붙이면 같은 Asset을 파티션 단위(예: 시간별)로 모델링할 수 있다. 파티션 키는 스케줄 주기로부터 미리 정해질 수도(사전 결정), 태스크 실행 중에만 정해질 수도(런타임) 있으며, 다운스트림에서는 PartitionedAssetTimetable과 파티션 매퍼로 업스트림 키를 다운스트림 키로 변환해 트리거 여부를 결정한다.

Asset 파티션이란

Asset 이벤트는 ``partition_key``를 가질 수 있으며, 이를 통해 같은 Asset을 파티션 단위로 세분화해서 모델링할 수 있다 (예: 시간별 파티션의 경우 ``2026-03-10T09:00:00``). Producer 쪽에서 매 실행마다 파티션 키가 있는 이벤트를 만들려면 ``CronPartitionTimetable``을 쓴다.

from airflow.sdk import CronPartitionTimetable, asset


@asset(
    uri="file://incoming/player-stats/team_b.csv",
    schedule=CronPartitionTimetable("15 * * * *", timezone="UTC"),
)
def team_b_player_stats():
    pass

파티션된 이벤트는 파티션을 인지하는 다운스트림 스케줄링 전용이며, 파티션을 인지하지 않는(non-partition-aware) 일반 Dag는 트리거하지 않는다.

핵심 포인트

  • partition_key는 같은 Asset을 파티션(예: 시간별) 단위로 세분화해서 모델링하는 값이며, CronPartitionTimetable로 매 실행마다 파티션 키가 있는 이벤트를 만들 수 있다
  • 파티션된 이벤트는 파티션을 인지하는 다운스트림 Dag만 트리거하고, 파티션을 인지하지 않는 일반 Dag는 트리거하지 않는다

사전 결정 파티셔닝 vs 런타임 파티셔닝

두 방식 모두 Dag 실행에 파티션 키를 붙이지만, '언제, 누가' 그 키를 정하느냐가 다르다. 사전 결정(pre-determined) 파티셔닝은 태스크가 실행되기 전에 타임테이블의 스케줄 주기와 파티션 매퍼로 업스트림 키를 다운스트림 키에 맞춰 미리 정한다 — producer 쪽에서는 ``CronPartitionTimetable``, consumer 쪽에서는 ``PartitionedAssetTimetable``이 이 방식을 쓴다. 런타임(runtime) 파티셔닝은 파티션 키 결정을 태스크 실행 시점까지 미룬다 — producing 태스크가 ``outlet_events[self].add_partitions(...)``로 키를 기록하며, ``PartitionedAtRuntime``이 이 방식을 쓰고 스스로는 절대 스케줄되지 않는다(``can_be_scheduled=False``). 커스텀 타임테이블도 ``CronTriggerTimetable``을 서브클래싱해 ``partitioned_at_runtime = True``를 설정하면 런타임 방식으로 파티션을 미룰 수 있다.

실무적인 선택 기준은 다음과 같다: 파티션 키가 스케줄 주기에서 자연스럽게 도출되면 ``CronPartitionTimetable``을, 소스 데이터에서 발견되는 워터마크처럼 태스크가 실행되어야만 키를 알 수 있으면 ``PartitionedAtRuntime``을 쓰고, 다운스트림에서는 어느 쪽이든 ``PartitionedAssetTimetable``로 소비한다. 한 타임테이블은 두 방식 중 하나만 쓰지, 둘을 섞어 쓰지 않는다.

핵심 포인트

  • 사전 결정 파티셔닝은 타임테이블 스케줄 주기+매퍼로 실행 전에 키가 정해지고(CronPartitionTimetable/PartitionedAssetTimetable), 런타임 파티셔닝은 태스크 실행 중 outlet_events[self].add_partitions(...)로 키가 정해진다(PartitionedAtRuntime, can_be_scheduled=False)
  • 파티션 키가 스케줄 주기에서 나오면 CronPartitionTimetable, 소스 데이터의 워터마크처럼 실행해야만 알 수 있으면 PartitionedAtRuntime을 쓰고, 다운스트림 소비는 어느 쪽이든 PartitionedAssetTimetable로 한다

PartitionedAssetTimetable과 기본 파티션 매퍼

``PartitionedAssetTimetable``은 파티션이 있는 이벤트만 소비한다 — partition_key가 없는 이벤트는 이 타임테이블을 쓰는 다운스트림 Dag를 트리거하지 않는다. ``default_partition_mapper``는 ``partition_mapper_config``로 개별 오버라이드하지 않는 한 모든 업스트림 Asset에 적용되며, 기본값은 키를 그대로 두는 ``IdentityMapper``다.

기본 제공 매퍼들의 동작은 다음과 같다.

* ``IdentityMapper`` — 키를 변형하지 않는다. * ``StartOfHourMapper``/``StartOfDayMapper``/``StartOfYearMapper`` — 시간 형태의 키를 지정한 단위로 정규화한다. 입력 키 ``2026-03-10T09:37:51``에 대해 각각 ``2026-03-10T09``, ``2026-03-10``, ``2026``을 출력한다. * ``ProductMapper`` — 복합 키를 세그먼트별로 매핑한 뒤 다시 결합한다. 예: 키 ``us|2026-03-10T09:00:00``에 ``ProductMapper(IdentityMapper(), StartOfDayMapper())``를 적용하면 ``us|2026-03-10``이 된다. * ``AllowedKeyMapper`` — 고정된 허용 목록에 있는 키만 통과시킨다. 예: ``AllowedKeyMapper(["us", "eu", "apac"])``는 이 세 지역 키만 허용하고 나머지는 거부한다. * ``FixedKeyMapper`` — 업스트림 키가 무엇이든 상관없이 고정된 다운스트림 키 하나로 접는다.

모든 필요한 업스트림 Asset의 변환된 키가 서로 일치하지 않으면 다운스트림 Dag는 트리거되지 않는다. 매퍼가 키를 변환할 수 없는 경우(예: 시간 형태를 기대하는 ``DailyMapper``에 ``"random-text"``가 들어온 경우)도 마찬가지로 트리거되지 않는다. 파티션된 Dag run 안에서는 ``dag_run.partition_key``로 해석된 파티션에 접근할 수 있고, 매퍼가 키를 datetime으로 해석할 수 있는 경우(``StartOf*Mapper`` 계열, IdentityMapper가 producer의 partition_date를 그대로 전달하는 경우, 그리고 유효 자식 매퍼가 시간형인 합성 매퍼)에는 ``dag_run.partition_date``로도 접근할 수 있다. ``ProductMapper``·``AllowedKeyMapper``나 ``to_partition_date``를 구현하지 않은 커스텀 매퍼는 결과 키가 날짜 형태처럼 보여도 partition_date를 ``None``으로 남겨두므로, 그런 소비자는 계속 partition_key를 직접 파싱해야 한다.

핵심 포인트

  • PartitionedAssetTimetable은 파티션 키가 없는 이벤트로는 트리거되지 않으며, default_partition_mapper의 기본값은 키를 그대로 두는 IdentityMapper다
  • StartOfHourMapper/StartOfDayMapper/StartOfYearMapper는 시간 키를 시/일/연 단위로 정규화하고, ProductMapper는 복합 키를 세그먼트별로 매핑하며, AllowedKeyMapper는 허용 목록 검증, FixedKeyMapper는 모든 키를 고정 키 하나로 접는다
  • 필요한 업스트림 키들이 서로 일치하지 않거나 매퍼가 키를 변환하지 못하면 다운스트림은 트리거되지 않으며, partition_date는 매퍼가 to_partition_date를 지원할 때만 채워진다