Asset 파티션 심화 (Rollup·FanOut 매퍼와 대기 정책)
Apache Airflow Official Documentation (in-repo snapshot) — Apache Software Foundation authoring-and-scheduling/assets.rst - Rollup mappers, Wait policies, Segment (categorical) rollup, Setting partition keys at runtime, Fan-out mappers, Window direction, Custom partition mappers/windows/timetables (versionadded 3.3.0)
이 모듈을 다 읽으면
- RollupMapper와 Window 조합이 여러 업스트림 파티션을 하나의 다운스트림 실행으로 묶는 원리를 설명할 수 있다
- WaitForAll과 MinimumCount 대기 정책의 차이를 판단할 수 있다
- SegmentWindow를 이용한 범주형(categorical) 롤업 구성 방법과 Window의 FORWARD/BACKWARD 방향·DST 주의사항을 설명할 수 있다
- FanOutMapper가 RollupMapper와 정반대로 동작하는 방식과 max_downstream_keys의 역할을 설명할 수 있다
RollupMapper는 여러 개의 세밀한 업스트림 파티션(예: 시간별)이 모두 도착해야 하나의 굵은 다운스트림 실행(예: 일별)을 트리거하도록 묶어주고, FanOutMapper는 반대로 하나의 업스트림 이벤트를 여러 다운스트림 실행으로 펼친다. 3.3.0에서 추가된 이 기능들은 Window·WaitPolicy·방향(direction) 조합으로 세밀하게 제어할 수 있으며, 플러그인으로 커스텀 매퍼·윈도·타임테이블도 확장할 수 있다.
RollupMapper: 여러 업스트림 파티션을 하나의 실행으로 묶기
시간별 업스트림이 일별 요약을 구동하거나, 일별 입력들이 주별 리포트를 구성하는 것처럼 더 굵은 다운스트림 기간이 여러 업스트림 이벤트로 이루어지는 경우 ``RollupMapper``를 쓴다. ``RollupMapper``는 각 업스트림 키를 다운스트림 그레인으로 정규화하는 upstream_mapper와, 하나의 다운스트림 키를 이루는 데 필요한 업스트림 키 전체 집합을 선언하는 ``Window``를 조합한다. 스케줄러는 윈도 안의 모든 업스트림 키가 도착할 때까지 Dag run을 보류하며, 부분적으로 채워진 윈도는 next-run-assets 화면에서 진행 상황을 확인할 수 있다.
기본 제공 윈도는 ``HourWindow``(60분), ``DayWindow``(24시간), ``WeekWindow``(7일), ``MonthWindow``, ``QuarterWindow``, ``YearWindow``이며, 같은 그레인으로 디코딩하는 upstream_mapper와 짝지어야 한다(예: ``StartOfHourMapper`` + ``DayWindow``).
RollupMapper(upstream_mapper=StartOfHourMapper(), window=DayWindow())
이렇게 하면 하루의 24개 시간별 파티션이 모두 도착해야 일별 요약 Dag가 한 번 실행된다. 잘못된 조합 — 예를 들어 문자열을 그대로 두는 identity 계열 매퍼를 ``DayWindow``(시간형 요구)와 짝짓는 경우 — 은 Dag 파싱 시점에 ``TypeError``를 발생시켜 설정 오류가 조용히 실행을 영원히 보류시키는 대신 즉시 드러나게 한다.
DST(서머타임) 관련 주의: ``DayWindow``는 항상 24개의 시간별 스텝을 나열한다. 서머타임을 관측하는 로컬 타임존으로 upstream_mapper를 구성하면, 봄철 시간이 앞당겨지는 날은 실제로 23시간뿐이라 윈도의 한 멤버가 영원히 매칭되지 않아 실행이 계속 보류되고, 가을철 시간이 되돌아가는 날은 25시간이라 중복된 시간이 버려진다. DST 경계를 넘는 롤업에는 UTC 기반 upstream_mapper를 쓰는 것이 권장된다.
핵심 포인트
- RollupMapper는 upstream_mapper(정규화)와 Window(필요한 업스트림 키 전체 집합)를 조합해, 윈도 안의 모든 업스트림 키가 도착해야 다운스트림 Dag run을 트리거한다
- 그레인이 맞지 않는 매퍼-윈도 조합(예: identity 매퍼 + DayWindow)은 Dag 파싱 시점에 TypeError로 즉시 실패해 무한 보류를 방지한다
- DayWindow는 항상 24시간을 요구하므로 DST가 있는 로컬 타임존을 쓰면 봄철(23시간, 영원히 보류)·가을철(25시간, 중복 시간 버려짐) 문제가 생기고, UTC 기반 매퍼가 권장된다
대기 정책: WaitForAll과 MinimumCount
``RollupMapper``의 ``wait_policy`` 인자는 기대되는 윈도 대비 실제로 도착한 업스트림 키를 기준으로 다운스트림 실행을 언제 발동할지 결정한다. 기본값인 ``WaitForAll``은 기대되는 모든 키가 도착할 때까지 실행을 보류한다. ``MinimumCount(n)``은 기대되는 키 중 최소 n개만 도착해도 조기에 발동하며, 느리거나 아예 오지 않는 업스트림 파티션을 허용해야 할 때 유용하다.
RollupMapper(
upstream_mapper=FixedKeyMapper("all_regions"),
window=SegmentWindow(["us", "eu", "apac"]),
wait_policy=MinimumCount(2),
)
이 예시는 3개 지역 파티션 중 2개만 도착해도 다운스트림이 발동한다. ``MinimumCount(-1)``은 '최대 1개까지 누락 허용'을 상대적으로 표현한 것으로, 3개짜리 윈도에서는 ``MinimumCount(2)``와 동일하다. 기본값을 그대로 쓰더라도 의도를 명시적으로 남기고 싶다면 ``WaitForAll()``을 직접 지정해도 된다.
핵심 포인트
- WaitForAll(기본값)은 윈도의 모든 키가 도착해야 발동하고, MinimumCount(n)은 n개만 도착해도 조기 발동해 느리거나 누락된 파티션을 허용한다
- MinimumCount(-1)은 '최대 1개 누락 허용'의 상대적 표현이며, 3개짜리 윈도에서 MinimumCount(2)와 동일하다
범주형(Segment) 롤업
지역·테넌트·실험 변형처럼 시간이 아닌 범주형 파티셔닝에는 ``SegmentWindow``와 ``FixedKeyMapper``를 조합한다. ``SegmentWindow(["us", "eu", "apac"])``는 하나의 다운스트림 기간을 이루는 고정된 문자열 키 집합을 선언하고, ``FixedKeyMapper("all_regions")``는 모든 업스트림 키를 단일 다운스트림 파티션 키로 접는다. 스케줄러는 선언된 모든 세그먼트가 도착할 때까지 다운스트림 Dag run을 보류했다가 한 번만 발동하며, 모든 세그먼트 이벤트는 하나의 ``AssetPartitionDagRun``에 누적되고 그 실행의 partition_key는 ``FixedKeyMapper``에 전달한 값이 된다. 이 조합은 ``WaitForAll``(기본값) 의미론에서만 의미가 있다.
구성 시점에 유효성 검증도 이루어진다: ``SegmentWindow``는 빈 리스트, 문자열이 아닌 항목, 빈 문자열 키에 대해 ``ValueError``를 던지며 중복 항목은 조용히 제거된다. ``FixedKeyMapper``는 인자가 비어있지 않은 문자열이 아니면 ``ValueError``를 던진다. 한 consumer Dag가 두 개 이상의 Asset을 롤업한다면 각 롤업마다 서로 다른 ``FixedKeyMapper`` 키를 줘서 (target_dag_id, partition_key) 버킷이 충돌하지 않게 해야 한다.
세그먼트 집합 자체를 런타임에 계산해야 하는 경우에는 여기(스케줄러 쪽)에 인코딩하지 말고, consumer 쪽 태스크에서 완결 여부를 직접 평가해야 한다 — 스케줄러는 파티션 집합을 결정하기 위해 사용자 코드를 실행해서는 안 되기 때문이다.
핵심 포인트
- SegmentWindow(고정 문자열 키 집합) + FixedKeyMapper(단일 다운스트림 키로 접기) 조합으로 범주형 롤업을 구성하며, 모든 세그먼트가 도착해야 한 번 발동한다(WaitForAll 의미론에서만 유효)
- SegmentWindow/FixedKeyMapper는 구성 시점에 입력을 검증해 ValueError를 던지고, 한 Dag가 여러 Asset을 롤업할 때는 서로 다른 FixedKeyMapper 키를 써서 버킷 충돌을 피해야 한다
런타임 파티션 키와 FanOutMapper
워터마크나 늦게 도착한 파일처럼 파티션 키가 실행 전에 정해지지 않는 경우, producer를 ``PartitionedAtRuntime()``으로 스케줄하고 ``outlet_events[self].add_partitions(...)``로 키를 기록한다. 단일 키 또는 리스트를 넘겨 한 실행에서 여러 파티션으로 팬아웃할 수 있고, 중복 키는 하나의 이벤트로 합쳐진다. 런타임 실행에서 정확히 하나의 파티션 키만 발행되면 그 실행의 ``dag_run.partition_key``도 그 값으로 소급 채워진다.
REST API로 파티션 키 기준 이벤트 조회도 가능하다: ``partition_key``는 정확 일치(B-tree 인덱스 사용, 항상 활성화), ``partition_key_regexp_pattern``은 정규식 필터링인데 ReDoS 위험 때문에 기본 비활성화되어 있고 ``[api] regexp_query_timeout``을 양수로 설정해야 켜진다(이 값이 쿼리 실행 시간도 제한한다). 정확한 키를 안다면 항상 활성화되어 있고 인덱스를 쓰는 ``partition_key``가 우선 선택지다.
``FanOutMapper``는 ``RollupMapper``의 정반대다 — 여러 업스트림 이벤트를 하나의 다운스트림 실행으로 모으는 대신, 하나의 업스트림 이벤트를 윈도 멤버 하나당 하나의 다운스트림 실행으로 펼친다. upstream_mapper(업스트림 키를 윈도 앵커로 정규화)와 Window(다운스트림 기간을 열거), 그리고 각 윈도 멤버를 다운스트림 파티션 키 문자열로 바꾸는 선택적 downstream_mapper로 구성된다. ``WeekWindow`` 같은 시간형 윈도는 기본 downstream_mapper(예: ``StartOfDayMapper``, ``YYYY-MM-DD`` 형식)가 자동 적용되지만, ``SegmentWindow``는 기본 매핑 테이블이 없어 downstream_mapper를 명시해야 한다.
FanOutMapper(upstream_mapper=StartOfWeekMapper(), window=WeekWindow(), max_downstream_keys=7)
이 예시는 주간 모델 아티팩트 하나를 요일별 추론 실행 7개로 펼친다. ``max_downstream_keys``는 하나의 업스트림 이벤트가 만들 수 있는 다운스트림 실행 개수의 상한이며, 초과하면 그 이벤트에 대한 실행들은 큐잉되지 않고 대신 'partition fan-out exceeded' 감사 로그 항목이 남는다. 생략하면 전역 설정 ``[scheduler] partition_mapper_max_downstream_keys``(기본 1000)를 따른다.
핵심 포인트
- 런타임 파티셔닝은 PartitionedAtRuntime() + outlet_events[self].add_partitions(...)로 구현하며, 정확 일치인 partition_key 필터는 항상 켜져 있고 인덱스를 쓰지만 정규식 필터는 ReDoS 우려로 기본 비활성화되어 있다
- FanOutMapper는 RollupMapper의 반대로, 하나의 업스트림 이벤트를 윈도 멤버 수만큼의 다운스트림 실행으로 펼치며, max_downstream_keys를 넘으면 큐잉되지 않고 감사 로그만 남는다(기본 상한은 전역 설정 1000)
Window 방향(FORWARD/BACKWARD)과 플러그인 확장
모든 ``Window``는 앵커를 기준으로 어느 기간을 열거할지 결정하는 ``direction`` 파라미터를 지원한다. 기본값 ``Window.Direction.FORWARD``는 업스트림 키에서 '시작하는' 기간을 산출한다 — 예를 들어 주간 키 ``"2026-03-09"``(월요일)에 대해 ``WeekWindow()``는 ``2026-03-09``부터 ``2026-03-15``까지 7일을 산출한다. ``Window.Direction.BACKWARD``는 그 키에서 '끝나는' 이전 기간을 산출한다 — 같은 키에 ``WeekWindow(direction=Window.Direction.BACKWARD)``를 쓰면 ``2026-03-03``부터 ``2026-03-09``까지가 산출된다. 이 방향 설정은 롤업 윈도에도 동일하게 적용되어, 예컨대 ``DayWindow(direction=Window.Direction.BACKWARD)``는 다운스트림 키 자정 '이전'의 24시간이 모두 도착할 때까지 실행을 보류한다.
커스텀 ``PartitionMapper``, ``Window``, 파티션 인지 ``Timetable``은 ``AirflowPlugin.partition_mappers``·``.windows``·``.timetables``에 등록해 플러그인으로 배포할 수 있다. 일단 플러그인이 설치되면 이 클래스들은 코어를 수정하지 않고도 ``PartitionedAssetTimetable``과 ``RollupMapper``에서 그대로 쓸 수 있다. 대표적인 확장 예시로는 네임스페이스 접두사를 제거하는 커스텀 파티션 매퍼(``"eu::daily-sales"``와 ``"us::daily-sales"``를 모두 ``"daily-sales"``로 접기), 평일만 열거하는 커스텀 롤업 윈도(주말 파티션을 기다리지 않도록), 그리고 파티션 키 결정을 태스크 런타임까지 미루는 커스텀 cron 타임테이블(그 기간의 데이터가 아직 준비되지 않았으면 파티션 키를 발행하지 않아, 파티션된 이벤트도 그 다운스트림 실행도 생기지 않게 하는 패턴)이 있다.
핵심 포인트
- Window.Direction.FORWARD(기본값)는 키에서 시작하는 기간을, BACKWARD는 키에서 끝나는 이전 기간을 산출하며 롤업 윈도에도 동일하게 적용된다
- 커스텀 PartitionMapper/Window/파티션 인지 Timetable은 AirflowPlugin에 등록해 플러그인으로 배포하면 코어 수정 없이 PartitionedAssetTimetable/RollupMapper에서 쓸 수 있다