Asset 정의

Asset 정의 (Asset Definitions)

이 페이지는 버전 2.4에 추가된 Asset(Airflow 3.0 이전에는 "Dataset"이라 불림)을 다뤄요. Asset은 데이터의 논리적 그룹이며, 상위 producer task가 asset을 업데이트하고, asset 업데이트가 하위 consumer DAG의 스케줄링에 기여해요. 유효한 URI란 무엇인지, asset에 여분 정보 붙이기, 보안 경고, asset 이벤트를 내보내는 task 만들기, 파티션, AssetAlias, 크로스 팀 필터링까지 다뤄요.

출처: 문서

본문

버전 2.4에 추가됨.

버전 3.0에서 변경: 이 개념은 이전에 "Dataset"이라 불렸어요.

"Asset"이란 무엇인가요?

Airflow asset은 데이터의 논리적 그룹이에요. 상위 producer task는 asset을 업데이트할 수 있고, asset 업데이트는 하위 consumer DAG의 스케줄링에 기여해요.

Uniform Resource Identifier (URI)가 asset을 정의해요:

from airflow.sdk import Asset

example_asset = Asset("s3://asset-bucket/example.csv")

Airflow는 URI가 나타내는 데이터의 내용이나 위치에 대해 어떤 가정도 하지 않으며 URI를 문자열처럼 취급해요. 즉 Airflow는 input_\d+.csv 같은 정규식이나 input_2022*.csv 같은 파일 glob 패턴을 하나의 선언에서 여러 asset을 만들려는 시도로 취급하며, 그것들은 작동하지 않아요.

유효한 URI로 asset을 만들어야 해요. Airflow core와 provider는 file(core), postgres(Postgres provider), s3(Amazon provider) 같은 사용할 수 있는 다양한 URI scheme을 정의해요. 서드파티 provider와 플러그인도 자체 scheme을 제공할 수 있어요. 이 미리 정의된 scheme들은 각각 따르기를 기대하는 개별 의미(semantics)가 있어요. 선택적 name 인자를 사용해 asset에 더 읽기 쉬운 식별자를 제공할 수 있어요.

from airflow.sdk import Asset

example_asset = Asset(uri="s3://asset-bucket/example.csv", name="bucket-1")

유효한 URI란?

기술적으로 URI는 RFC 3986의 유효 문자 집합, 즉 기본적으로 ASCII 영숫자 문자와 %, -, _, ., ~를 따라야 해요. URI 안전 문자로 표현할 수 없는 리소스를 식별하려면 percent-encoding으로 리소스 이름을 인코딩해요.

URI는 또한 대소문자를 구분하므로, s3://example/assets3://Example/asset은 다른 것으로 간주돼요. URI의 host 부분도 대소문자를 구분한다는 점에 유의하세요. 이는 RFC 3986과 다릅니다.

미리 정의된 scheme(예: file, postgres, s3)의 경우 의미 있는 URI를 제공해야 해요. 제공할 수 없다면, 의미 제한이 없는 다른 scheme을 사용해요. Airflow는 사용자 정의 URI scheme(x- 접두사)에 대해 결코 의미를 요구하지 않으므로, 그것이 좋은 대안이 될 수 있어요. 나중에만 얻을 수 있는 URI(예: Task 실행 중)가 있다면 AssetAlias 사용을 고려하고 URI를 나중에 업데이트해요.

# invalid asset:
must_contain_bucket_name = Asset("s3://")

Airflow 내부용으로 예약된 airflow scheme을 사용하지 마세요.

Airflow는 scheme에서 항상 소문자를 선호하며, 리소스를 올바르게 구별하려면 URI의 host 부분에 대소문자 구분이 필요해요.

# invalid assets:
reserved = Asset("airflow://example_asset")
not_ascii = Asset("èxample_datašet")

추가 의미 제약이 없는 scheme으로 asset을 정의하려면 x- 접두사가 있는 scheme을 사용해요. Airflow는 이 scheme의 URI에 대한 의미 검증을 건너뛰어요.

# valid asset, treated as a plain string
my_ds = Asset("x-my-thing://foobarbaz")

식별자는 절대적일 필요가 없어요. scheme이 없는 상대 URI이거나 단순한 경로·문자열일 수도 있어요:

# valid assets:
schemeless = Asset("//example/asset")
csv_file = Asset("example_asset")

비절대 식별자는 Airflow에 아무 의미도 전달하지 않는 일반 문자열로 간주돼요.

Asset에 대한 추가 정보

필요하면 extra 파라미터로 asset에 추가 dict를 포함할 수 있어요:

example_asset = Asset(
    "s3://asset/example.csv",
    extra={"team": "trainees"},
)

이를 통해 소유권 정보나 파일의 목적 같은 asset에 대한 커스텀 메타데이터를 제공할 수 있어요. extra 필드는 asset의 정체성에 영향을 주지 않아요. 따라서 extra 값의 고유성을 유지하는 것은 사용자 책임이에요. asset당 단일 extra 값 집합만 두는 것이 좋아요.

예를 들어 다음 스니펫에서는 extra dict 중 하나만 결국 저장되지만, 어느 것이 저장될지는 보장되지 않아요.

Asset("s3://asset/example.csv", extra={"d": "e"})
Asset("s3://asset/example.csv", extra={"f": "g"})

이 동작은 AssetAlias로 생성된 동적 생성 asset에도 적용돼요. 아래 예제에서 최종 저장된 extra 값은 보장되지 않으며 Dag processor 설정에 따라 달라질 수 있어요.

from airflow.sdk import AssetAlias

@dag(schedule=None)
def my_dag_1():

    @task(outlets=[AssetAlias("my-task-outputs")])
    def my_task_with_outlet_events(*, outlet_events):
        outlet_events[AssetAlias("my-task-outputs")].add(
            # Asset extra set as {"from": "asset alias"}
            Asset("s3://bucket/my-task", extra={"from": "asset alias"})
        )

    my_task_with_outlet_events()

# Asset extra set as {"key": "value"}
@dag(schedule=Asset("s3://bucket/my-task", extra={"key": "value"}))
def my_dag_2(): ...

my_dag_1()
my_dag_2()

# It's not guaranteed which extra will be the one stored

보안 경고

  1. Asset URI의 안전한 명명: Asset URI와 extra 필드의 값은 Airflow 메타데이터 데이터베이스에 평문(cleartext)으로 저장돼요. 이 필드는 암호화되지 않아요. 민감한 정보, 특히 자격 증명을 asset URI나 extra dict에 절대 저장하지 마세요.
  2. Asset 생성의 보안 함의: Airflow 보안 모델에서 Assets에 대한 can_create 권한을 부여하는 것은, 그 asset에 의존하는 모든 하위 DAG에 대해 "trigger" 권한을 부여하는 것과 실질적으로 동일해요. Airflow는 데이터 인지 스케줄링에 "암묵적 신뢰(implicit trust)" 모델을 사용하므로, Asset Event를 생성할 수 있는 사용자(API나 task를 통해)는 하위 DAG를 조회·편집할 권한이 없어도 그 asset에 스케줄링된 어떤 Dag든 트리거할 수 있어요. 멀티 테넌트 환경에서 Assets에 can_create를 부여할 때 주의하세요. 사용자가 자신의 직접 범위 밖의 워크플로에 영향을 줄 수 있게 하니까요.

Asset 이벤트를 내보내는 Task 만들기

asset이 정의되면 Task가 outlets를 지정해 그것에 대한 이벤트를 내보낼 수 있어요:

from airflow.sdk import DAG, Asset
from airflow.providers.standard.operators.python import PythonOperator

example_asset = Asset(name="example_asset", uri="s3://asset-bucket/example.csv")

def _write_example_asset():
    """Write data to example_asset..."""

with DAG(dag_id="example_asset", schedule="@daily"):
    PythonOperator(task_id="example_asset", outlets=[example_asset], python_callable=_write_example_asset)

이것은 꽤 많은 보일러플레이트예요. Airflow는 하나의 asset의 이벤트를 내보내는 단일 task를 가진 DAG를 만드는 이 단순하지만 가장 흔한 경우에 대한 축약을 제공해요. 아래 코드 블록은 위와 정확히 동등해요:

from airflow.sdk import asset

@asset(uri="s3://asset-bucket/example.csv", schedule="@daily")
def example_asset():
    """Write data to example_asset..."""

@asset을 선언하면 자동으로 다음이 생성돼요:

  • name이 함수 이름으로 설정된 Asset.
  • dag_id가 함수 이름으로 설정된 DAG.
  • task_id가 함수 이름으로 설정되고 outlet이 생성된 Asset인, DAG 안의 task.

self, context, outlet_events 파라미터 이름은 @asset 함수에서 예약되어 있어요: Airflow가 런타임에 (각각 asset 자체, 실행 컨텍스트, outlet 이벤트 접근자로) 채우며, 절대 inlet asset 참조로 취급되지 않아요.

내보내는 asset 이벤트에 추가 정보 붙이기

버전 2.10.0에 추가됨.

asset outlet이 있는 task는 asset 이벤트를 내보내기 전에 선택적으로 추가 정보를 붙일 수 있어요. 이는 Asset에 대한 추가 정보와 다릅니다. asset에 대한 추가 정보는 asset URI가 가리키는 엔티티를 정적으로 설명해요. asset 이벤트의 추가 정보는 대신 트리거하는 데이터 변경을 주석으로 다는 데 사용해야 해요 — 예: 업데이트가 데이터베이스에서 변경한 행 수, 또는 그것이 덮는 날짜 범위처럼요.

asset 이벤트에 추가 정보를 붙이는 가장 쉬운 방법은 task에서 Metadata 객체를 yield하는 것이에요:

from airflow.sdk import Metadata, asset

@asset(uri="s3://asset/example.csv", schedule=None)
def example_s3(self):  # 'self' here refers to the current asset.
    df = ...  # Get a Pandas DataFrame to write.
    # Write df to asset...
    yield Metadata(self, {"row_count": len(df)})

Airflow는 yield된 모든 메타데이터를 자동으로 수집하고, 해당 metadata 객체에 대한 추가 정보로 asset 이벤트를 채워요.

이것은 classic operator에서도 할 수 있어요. 가장 좋은 방법은 operator를 서브클래싱하고 execute를 오버라이드하는 것이에요. 또는 task의 pre_executepost_execute 훅에서 extras를 추가할 수 있어요. 훅을 사용한다면, 훅은 task가 재시도될 때 다시 실행되지 않으며 어떤 시나리오에서는 추가 정보가 실제 데이터와 일치하지 않을 수 있다는 점을 기억하세요.

같은 것을 달성하는 또 다른 방법은 task의 실행 컨텍스트에서 outlet_events에 직접 접근하는 것이에요:

@asset(schedule=None)
def write_to_s3(self, context):
    context["outlet_events"][self].extra = {"row_count": len(df)}

여기엔 마법이 거의 없어요 — Airflow는 단순히 yield된 값을 정확히 같은 접근자에 써요. 이것도 execute, pre_execute, post_execute를 포함한 classic operator에서 작동해요.

Note

Asset 이벤트 추가 정보는 JSON 직렬화 가능한 값만 포함할 수 있어요(list와 dict 중첩 가능). 이는 값이 데이터베이스에 저장되기 때문이에요.

이전에 내보낸 asset 이벤트에서 정보 가져오기

버전 2.10.0에 추가됨.

이전 섹션에서 설명한 대로 task의 outlets에 정의된 asset의 이벤트는, inlets에 같은 asset을 선언하는 task가 읽을 수 있어요. asset 이벤트 항목은 extra(이전 섹션 참고), task에서 이벤트가 내보내진 시점을 나타내는 timestamp, 이벤트를 소스로 연결하는 source_task_instance를 포함해요.

Inlet asset 이벤트는 실행 컨텍스트의 inlet_events 접근자로 읽을 수 있어요. 이전 섹션의 write_to_s3 asset에서 계속하면:

@asset(schedule=None)
def post_process_s3_file(context, write_to_s3):  # Declaring an inlet to write_to_s3.
    events = context["inlet_events"][write_to_s3]
    last_row_count = events[-1].extra["row_count"]

inlet_events 매핑의 각 값은 주어진 asset의 과거 이벤트를 timestamp별로, 가장 이른 것부터 최신 것 순으로 정렬하는 시퀀스류 객체예요. 대부분의 Python list 인터페이스를 지원하므로, [-1]로 마지막 이벤트에, [-2:]로 마지막 두 개에 접근하는 등 할 수 있어요. 접근자는 lazy이며 항목에 접근할 때만 데이터베이스에 접속해요.

@asset, @task, classic operator 사이의 의존성

@asset은 task와 asset을 가진 DAG의 래퍼일 뿐이므로, @task나 classic operator에서 @asset을 읽기는 매우 쉽습니다. 예를 들어 위의 post_process_s3_file은 task로도 작성할 수 있어요(간결성을 위해 DAG 안쪽은 생략):

@task(inlets=[write_to_s3])
def post_process_s3_file(*, inlet_events):
    events = inlet_events[example_s3_asset]
    last_row_count = events[-1].extra["row_count"]

post_process_s3_file()

반대도 적용돼요:

example_asset = Asset("example_asset")

@task(outlets=[example_asset])
def emit_example_asset():
    """Write to example_asset..."""

@asset(schedule=None)
def process_example_asset(example_asset):
    """Process inlet example_asset..."""

또한 @asset@task와 함께 사용해 asset을 생성하는 task를 커스터마이즈할 수 있어요. TaskFlow API로 Pythonic Dags에 설명된 현대적 TaskFlow 접근을 활용해요.

이 조합을 통해 task의 초기 인자를 설정하고 BashOperator 같은 다양한 operator를 사용할 수 있어요:

@asset(schedule=None)
@task.bash(retries=3)
def example_asset():
    """Write to example_asset, from a Bash task with 3 retries..."""
    return "echo 'run'"

하나의 Task에서 여러 asset으로 출력

task가 여러 asset에 대한 이벤트를 내보내는 것이 가능해요. 일반적으로 권장되지 않지만, 데이터 소스를 여러 개로 나눠야 하는 등 특정 상황에서 필요해요. outlets가 설계상 복수형이므로 task에서 이는 간단해요:

from airflow.sdk import DAG, Asset, task

input_asset = Asset("input_asset")
out_asset_1 = Asset("out_asset_1")
out_asset_2 = Asset("out_asset_2")

with DAG(dag_id="process_input", schedule=None):

    @task(inlets=[input_asset], outlets=[out_asset_1, out_asset_2])
    def process_input():
        """Split input into two."""

이에 대한 축약은 @asset.multi예요:

from airflow.sdk import Asset, asset

input_asset = Asset("input_asset")
out_asset_1 = Asset("out_asset_1")
out_asset_2 = Asset("out_asset_2")

@asset.multi(schedule=None, outlets=[out_asset_1, out_asset_2])
def process_input(input_asset):
    """Split input into two."""

AssetAlias를 통한 동적 데이터 이벤트 발행·asset 생성

task가 Asset의 고정 속성(URI나 name 같은)을 알기 전에 asset 의존성을 선언해야 한다면 AssetAlias를 사용해요. 별칭(alias)은 안정적인 이름으로 outlets에 나열되고, task는 outlet_events나 yield된 Metadata로 하나 이상의 구체적 Asset 객체를 추가해 런타임에 해석해요. 하위 DAG는 별칭에 의존할 수 있고, Airflow는 해결된 asset에 대해 이벤트가 내보내질 때 그것들을 트리거해요.

AssetAlias 사용법

AssetAlias는 별칭을 고유하게 식별하는 단일 인자 name을 가져요. Task는 먼저 별칭을 outlet으로 선언한 다음, 실행 중 outlet_events를 사용하거나 Metadata를 yield해 별칭을 그것이 생성한 구체적 asset과 연결해야 해요.

다음 예제는 선택적 추가 정보 extra와 함께 S3 URI f"s3://bucket/my-task"에 대해 asset 이벤트를 생성해요. asset이 없으면 Airflow가 동적으로 생성하고 경고 메시지를 기록해요.

outlet_events를 통해 task 실행 중 asset 이벤트 내보내기

from airflow.sdk import AssetAlias

@task(outlets=[AssetAlias("my-task-outputs")])
def my_task_with_outlet_events(*, outlet_events):
    outlet_events[AssetAlias("my-task-outputs")].add(Asset("s3://bucket/my-task"), extra={"k": "v"})

Metadata를 yield해 task 실행 중 asset 이벤트 내보내기

from airflow.sdk import Metadata

@task(outlets=[AssetAlias("my-task-outputs")])
def my_task_with_metadata():
    s3_asset = Asset(uri="s3://bucket/my-task", name="example_s3")
    yield Metadata(s3_asset, extra={"k": "v"}, alias=AssetAlias("my-task-outputs"))

추가된 asset에 대해 하나의 asset 이벤트만 내보내져요. 별칭에 여러 번 추가되거나 여러 별칭에 추가돼도요. 하지만 다른 extra 값이 전달되면 여러 asset 이벤트를 내보낼 수 있어요. 다음 예제에서는 두 개의 asset 이벤트가 내보내질 거예요.

from airflow.sdk import AssetAlias

@task(
    outlets=[
        AssetAlias("my-task-outputs-1"),
        AssetAlias("my-task-outputs-2"),
        AssetAlias("my-task-outputs-3"),
    ]
)
def my_task_with_outlet_events(*, outlet_events):
    outlet_events[AssetAlias("my-task-outputs-1")].add(Asset("s3://bucket/my-task"), extra={"k": "v"})
    # This line won't emit an additional asset event as the asset and extra are the same as the previous line.
    outlet_events[AssetAlias("my-task-outputs-2")].add(Asset("s3://bucket/my-task"), extra={"k": "v"})
    # This line will emit an additional asset event as the extra is different.
    outlet_events[AssetAlias("my-task-outputs-3")].add(Asset("s3://bucket/my-task"), extra={"k2": "v2"})

해결된 asset 별칭을 통해 이전에 내보낸 asset 이벤트에서 정보 가져오기

이전에 내보낸 asset 이벤트에서 정보 가져오기에서 언급했듯이 inlet asset 이벤트는 실행 컨텍스트의 inlet_events 접근자로 읽을 수 있으며, asset aliases를 사용해 그것들이 트리거한 asset 이벤트에 접근할 수도 있어요.

with DAG(dag_id="asset-alias-producer"):

    @task(outlets=[AssetAlias("example-alias")])
    def produce_asset_events(*, outlet_events):
        outlet_events[AssetAlias("example-alias")].add(Asset("s3://bucket/my-task"), extra={"row_count": 1})

with DAG(dag_id="asset-alias-consumer", schedule=None):

    @task(inlets=[AssetAlias("example-alias")])
    def consume_asset_alias_events(*, inlet_events):
        events = inlet_events[AssetAlias("example-alias")]
        last_row_count = events[-1].extra["row_count"]

producer_teams로 크로스 팀 asset 이벤트 필터링

버전 3.3.0에 추가됨.

Multi-Team 모드가 활성화되면 asset 이벤트가 팀 구성원으로 필터링돼요. 기본적으로 consumer DAG는 같은 팀 내의 DAG나 global(팀 없는) DAG가 생성한 asset 이벤트만 받아요. 이는 의도하지 않은 크로스 팀 트리거를 방지해요.

크로스 팀 접근을 구성하려면 Asset 정의의 access_control 파라미터에 AssetAccessControl 인스턴스를 사용해요:

from airflow.sdk import Asset, AssetAccessControl

shared_data = Asset(
    name="my_data",
    uri="s3://bucket/shared/data.csv",
    access_control=AssetAccessControl(
        producer_teams=["team_analytics", "team_ml"],
    ),
)

이 예제에서 team_analyticsteam_ml에 속한 DAG가 생성한 asset 이벤트는, consumer DAG 자신의 팀의 이벤트에 더해, shared_data에 스케줄링하는 모든 consumer DAG가 수락해요.

AssetAccessControl 파라미터

AssetAccessControl 클래스는 다음 파라미터를 받아요:

  • producer_teams (list[str], 기본 []): 이 asset의 consumer가 소비하는 이벤트를 생성할 수 있는 팀 이름 목록. consumer 자신의 팀에 더해요.
  • consumer_teams (list[str] | None, 기본 None): 이 asset의 producer가 생성한 이벤트를 소비할 수 있는 팀 이름 목록. consumer_teams로 크로스 팀 asset 이벤트 필터링 참고.
  • allow_global (bool, 기본 True): 팀 없는(global) DAG가 크로스 팀 이벤트 전달에 참여할 수 있는지 여부. consumer측과 producer측 asset 모두에 대한 전체 의미는 Multi-Team을 참고해요.

Global producer 차단

기본적으로 global(팀 없는) DAG는 어떤 consumer든 트리거할 수 있어요. 엄격한 팀 격리 시나리오에서는 팀 없는 producer를 차단하고 싶을 수 있어요:

from airflow.sdk import Asset, AssetAccessControl

strict_data = Asset(
    name="strict_data",
    uri="s3://bucket/strict/data.csv",
    access_control=AssetAccessControl(
        producer_teams=["team_analytics"],
        allow_global=False,
    ),
)

allow_global=False면 consumer 자신의 팀이나 team_analytics에 속한 DAG만 strict_data의 consumer를 트리거할 수 있어요. 팀 없는 Dag producer는 차단돼요.

Note

allow_global 플래그는 Dag producer에만 영향을 줘요. 팀 없는 API 사용자는 이 설정과 무관하게 항상 팀 없는 consumer만 트리거하도록 제한돼요.

기본 동작

access_control을 지정하지 않으면 기본 AssetAccessControl()이 사용돼요(빈 producer_teams, consumer_teams=None, allow_global=True). 완전한 동작 규칙 표는 Multi-Team을 참고해요. 요약하면, 규칙은 producer와 consumer가 팀 연관을 가지는지에 따라 달라져요:

  • 둘 다 같은 팀: 이벤트는 항상 전달돼요.
  • Producer가 팀, consumer는 다른 팀: 이벤트는 차단돼요(producer의 팀이 asset의 producer_teams에 없으면).
  • Producer가 팀 없음 (global Dag): 이벤트는 asset의 allow_global=True(기본값)인 모든 consumer에게 전달돼요. Global DAG는 어떤 팀이든 의존할 수 있는 공유 인프라 역할을 해요.
  • Consumer가 팀 없음 (global Dag): consumer는 producer의 팀과 무관하게 어느 소스에서든 이벤트를 수락해요. 팀 없는 consumer는 어떤 팀이든 공급할 수 있는 공유 인프라 역할을 해요.
  • 둘 다 팀 없음: 이벤트는 전달돼요 (둘 다 global).

Multi-Team 모드가 비활성화되면 access_control은 무시되고 모든 asset 이벤트가 모든 consumer DAG로 전달되어 하위 호환성을 보존해요.

Asset 파티션

버전 3.2.0에 추가됨.

Asset 이벤트는 partition_key를 포함해 partitioned(파티션화)될 수 있어요. 이를 통해 같은 asset을 파티션 세분성으로 모델링할 수 있어요(예: 시간별 파티션의 경우 2026-03-10T09:00:00).

스케줄에 따라 파티션 이벤트를 생성하려면 producer Dag(또는 @asset)에서 CronPartitionTimetable을 사용해요. 이 timetable은 각 run에서 partition key를 가진 asset 이벤트를 생성해요.

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

파티션 이벤트는 파티션 인지 하위 스케줄링을 위한 것이며, 파티션 인지가 아닌 DAG는 트리거하지 않아요.

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

두 종류 모두 Dag run에 파티션 키를 붙여요 — 차이는 키가 언제, 누구에 의해 결정되느냐예요:

  • 사전 결정 파티셔닝 — 파티션 키가 task 실행 전에 결정돼요. timetable의 스케줄 주기와 파티션 매퍼를 사용해 상위 키를 하위 키와 일치시키고 파티션 기반 Dag run을 트리거해요. CronPartitionTimetable은 producer로 이것을 사용하고, PartitionedAssetTimetable은 consumer로 이것을 사용해요.
  • 런타임 파티셔닝 — 파티션 키가 task 런타임으로 연기돼요: producer task가 outlet_events[self].add_partitions(...)로 키를 기록해요. PartitionedAtRuntime은 이 종류를 사용하며 스스로 스케줄링하지 않아요(can_be_scheduled=False). 스케줄링 가능한 timetable도 CronTriggerTimetable을 서브클래싱하고 partitioned_at_runtime = True를 설정해 런타임으로 연기할 수 있어요(아래 커스텀 플러그인 예제 참고).

timetable은 한 종류 또는 다른 종류를 사용하지, 둘 다는 아니에요: run 전에 파티션을 해결하거나 task 런타임으로 연기해요.

실용 규칙: 파티션 키가 스케줄 주기에서 나올 때는 CronPartitionTimetable을, 키가 task가 실행된 후에만 알려질 때(예: 소스 데이터의 watermark)는 PartitionedAtRuntime을, 하위에서 어느 종류든 소비하려면 PartitionedAssetTimetable을 사용해요.

하위 파티션 인지 스케줄링에는 PartitionedAssetTimetable을 사용해요:

from airflow.sdk import DAG, StartOfHourMapper, PartitionedAssetTimetable

with DAG(
    dag_id="clean_and_combine_player_stats",
    schedule=PartitionedAssetTimetable(
        assets=team_a_player_stats & team_b_player_stats & team_c_player_stats,
        default_partition_mapper=StartOfHourMapper(),
    ),
    catchup=False,
):
    ...

PartitionedAssetTimetable은 파티션화된 asset 이벤트를 요구해요. asset 이벤트에 partition_key가 없으면, PartitionedAssetTimetable을 사용하는 하위 DAG를 트리거하지 않아요.

default_partition_mapperpartition_mapper_config로 오버라이드하지 않는 한 모든 상위 asset에 사용돼요. 기본 매퍼는 IdentityMapper(키 변환 없음)예요.

파티션 매퍼는 상위 파티션 키가 하위 DAG 파티션 키로 어떻게 변환되는지 정의해요:

  • IdentityMapper는 키를 변경하지 않고 유지해요.
  • StartOfHourMapper, StartOfDayMapper, StartOfYearMapper 같은 시간 매퍼는 시간 키를 선택된 세분성으로 정규화해요. 입력 키 2026-03-10T09:37:51에 대한 기본 출력은:
    • StartOfHourMapper -> 2026-03-10T09
    • StartOfDayMapper -> 2026-03-10
    • StartOfYearMapper -> 2026
  • ProductMapper는 복합 키를 세그먼트별로 매핑해요. 세그먼트당 하나의 매퍼를 적용한 다음 매핑된 세그먼트를 다시 결합해요. 예를 들어 키 us|2026-03-10T09:00:00ProductMapper(IdentityMapper(), StartOfDayMapper())us|2026-03-10을 만듭니다.
  • AllowedKeyMapper는 키가 고정 allow-list에 있는지 검증하고, 유효하면 키를 변경 없이 통과시켜요. 예를 들어 AllowedKeyMapper(["us", "eu", "apac"])는 그 지역 키들만 수락하고 나머지를 거부해요.
  • FixedKeyMapper는 상위 값과 무관하게 모든 상위 키를 고정된 하위 키로 붕괴시켜요.
  • SegmentWindow는 하나의 하위 기간을 구성하는 고정 범주형 문자열 키 집합(예: 지역, 테넌트)을 선언해요. RollupMapper 안에서 FixedKeyMapper와 짝을 이루면, 선언된 모든 세그먼트가 도착할 때까지 하위 run을 보유해요(segment-rollup 참고).

Per-asset 매퍼 구성과 복합 키 매핑 예제:

from airflow.sdk import (
    Asset,
    IdentityMapper,
    PartitionedAssetTimetable,
    ProductMapper,
    StartOfDayMapper,
)

regional_sales = Asset(uri="file://incoming/sales/regional.csv", name="regional_sales")

with DAG(
    dag_id="aggregate_regional_sales",
    schedule=PartitionedAssetTimetable(
        assets=regional_sales,
        default_partition_mapper=ProductMapper(IdentityMapper(), StartOfDayMapper()),
    ),
):
    ...

특정 상위 asset에 대한 매퍼는 partition_mapper_config로 오버라이드할 수도 있어요:

from airflow.sdk import Asset, DAG, StartOfDayMapper, IdentityMapper, PartitionedAssetTimetable

hourly_sales = Asset(uri="file://incoming/sales/hourly.csv", name="hourly_sales")
daily_targets = Asset(uri="file://incoming/sales/targets.csv", name="daily_targets")

with DAG(
    dag_id="join_sales_and_targets",
    schedule=PartitionedAssetTimetable(
        assets=hourly_sales & daily_targets,
        # Default behavior: map timestamp-like keys to daily keys.
        default_partition_mapper=StartOfDayMapper(),
        # Override for assets that already emit daily partition keys.
        partition_mapper_config={
            daily_targets: IdentityMapper(),
        },
    ),
):
    ...

모든 필수 상위 asset의 변환된 파티션 키가 정렬되지 않으면, 그 파티션에 대해 하위 DAG가 트리거되지 않아요.

매퍼가 키를 변환할 수 없을 때도 동일하게 적용돼요. 예를 들어 상위 이벤트가 partition_key="random-text"이고 하위 매핑이 DailyMapper(시간 유사 키 기대)를 사용하면, 하위 파티션 일치가 생성될 수 없으므로 그 키에 대해 하위 DAG가 트리거되지 않아요.

파티션된 Dag run 안에서 dag_run.partition_key로 해결된 파티션에 접근해요. consumer의 파티션 매퍼가 키를 datetime으로 해결할 수 있으면 그 값은 dag_run.partition_date로도 사용할 수 있어서, 템플릿은 {{ partition_date | ds }}를 사용할 수 있어요. 이는 StartOf*Mapper 계열(키를 직접 디코드), IdentityMapper(producer의 partition_date를 통과), 그리고 RollupMapper, ChainMapper, FanOutMapper 같은 복합 매퍼 — 그 유효 child 매퍼가 시간형인 것 — 을 다룹니다. 키에 시간 의미가 없는 매퍼(ProductMapper, AllowedKeyMapperto_partition_date를 구현하지 않는 커스텀 매퍼)는 결과 키가 날짜 모양이어도 partition_dateNone으로 두므로, 그런 consumer는 partition_key 파싱을 계속해야 해요.

DagRun을 파티션 키로 수동 트리거할 수도 있어요 (예: UI의 Trigger Dag 창, 또는 REST API로 요청 본문에 partition_key 포함):

curl -X POST "http://<airflow-host>/api/v2/dags/aggregate_regional_sales/dagRuns" \
  -H "Content-Type: application/json" \
  -d '{
    "logical_date": "2026-03-10T00:00:00Z",
    "partition_key": "us|2026-03-10T09:00:00"
  }'

롤업 매퍼

버전 3.3.0에 추가됨.

위에 보여준 매퍼들은 상위 키를 단일 하위 키에 1:1로 일치시켜요. 많은 상위 이벤트로 구성된 더 거친 하위 기간 — 하루 요약을 구동하는 시간별 상위, 주간 보고서를 구성하는 일별 입력 — 에는 RollupMapper를 사용해요. RollupMapper는 상위 매퍼(각 상위 키를 하위 세분성으로 정규화)와, 하나의 하위 키에 필요한 전체 상위 키 집합을 선언하는 Window를 구성해요. Scheduler는 window의 모든 상위 키가 도착할 때까지 Dag run을 보유하고, 부분 windows는 next-run-assets 뷰에 pending으로 남아 operator들이 진행을 볼 수 있어요.

제공되는 windows는 HourWindow(시간당 60분), DayWindow(하루당 24시간), WeekWindow(주당 7일), MonthWindow, QuarterWindow, YearWindow이에요. 각 window를 같은 시간 세분성으로 디코드하는 상위 매퍼와 짝지어요 — 예를 들어 StartOfHourMapperDayWindow와 함께요.

각 상위 asset 이벤트는 2026-03-10T09:00:00(초 정밀도) 같은 미세한 파티션 키를 운반해요. StartOfHourMapper는 그 키를 시간 경계 2026-03-10T09로 정규화하는데, 이는 기대되는 각 멤버를 인코딩하는 형식이에요 — scheduler가 DayWindow의 스물네 개 필수 멤버와 대조하는 것과 같은 문자열이에요.

다음 시간별-일별 예제는 달력 날짜의 스물네 개 시간별 상위 파티션이 모두 도착하면 일일 요약을 생성해요:

from airflow.sdk import (
    DAG,
    Asset,
    CronPartitionTimetable,
    DayWindow,
    PartitionedAssetTimetable,
    RollupMapper,
    StartOfHourMapper,
    task,
)

hourly_sales = Asset(uri="file://incoming/sales/hourly.csv", name="hourly_sales")

# Producer: emits one partitioned event per hour (key looks like 2026-03-10T09:00:00).
with DAG(
    dag_id="ingest_hourly_sales",
    schedule=CronPartitionTimetable("0 * * * *", timezone="UTC"),
):

    @task(outlets=[hourly_sales])
    def ingest():
        pass

    ingest()

# Consumer: fires once a day's twenty-four hourly partitions are all in.
with DAG(
    dag_id="daily_sales_summary",
    schedule=PartitionedAssetTimetable(
        assets=hourly_sales,
        default_partition_mapper=RollupMapper(
            upstream_mapper=StartOfHourMapper(),
            window=DayWindow(),
        ),
    ),
    catchup=False,
):

    @task
    def summarize(dag_run=None):
        # dag_run.partition_key is the day, e.g. "2026-03-10".
        print(dag_run.partition_key)

    summarize()

잘못 구성된 RollupMapper — 예를 들어 identity 디코딩 상위 매퍼를 DayWindow와 짝지은 — 는 Dag 파싱 시 TypeError를 발생시켜 오작동이 하위 모든 run을 조용히 붙잡기 전에 즉시 표면화돼요. 에러는 타입 불일치예요: identity 디코딩 매퍼의 expected_decoded_typestr이지만, DayWindow 같은 시간 windows는 datetime을 요구해요. RollupMapper는 생성 시점에 불일치를 감지하고 Dag가 스케줄링되기 전에 발생시켜요.

DayWindow는 항상 스물네 개의 시간별 단계를 열거해요. 일광 절약 시간(DST)을 관찰하는 로컬 타임존으로 구성된 상위 매퍼와 함께라면, spring-forward 날에는 실제 시간이 스물세 시간뿐이라(한 window 멤버에게 일치 이벤트가 절대 없어 run이 무기한 보유됨) fall-back 날에는 스물다섯 시간(반복된 시간이 버려짐)이 돼요. DST 경계를 가로지르는 롤업에는 UTC 기반 상위 매퍼를 사용해요. 전체 논의는 DayWindow 클래스 docstring을 참고해요.

대기 정책

RollupMapper는 선택적 wait_policy 인자를 받는데, 기대되는 window와 실제 도착한 상위 키가 주어졌을 때 하위 Dag run이 언제 발화할지 결정해요.

  • WaitForAll(기본값)은 window의 모든 기대되는 상위 키가 도착할 때까지 run을 보유해요.
  • MinimumCount (n)은 기대되는 키 중 적어도 n개가 도착하면 일찍 발화해요 — 느리거나 누락된 상위 파티션을 run을 무기한 보유하기보다는 허용하는 데 유용해요.
from airflow.sdk import (
    DAG,
    Asset,
    FixedKeyMapper,
    MinimumCount,
    PartitionedAtRuntime,
    PartitionedAssetTimetable,
    RollupMapper,
    SegmentWindow,
    asset,
)

@asset(
    uri="file://incoming/player-stats/multi-region.csv",
    schedule=PartitionedAtRuntime(),
)
def multi_region_player_stats(self, outlet_events):
    outlet_events[self].add_partitions(["us", "eu", "apac"])

# Consumer: fires once at least two of the three declared region partitions arrive.
with DAG(
    dag_id="segment_region_stats_early_rollup",
    schedule=PartitionedAssetTimetable(
        assets=Asset.ref(name="multi_region_player_stats"),
        default_partition_mapper=RollupMapper(
            upstream_mapper=FixedKeyMapper("all_regions"),
            window=SegmentWindow(["us", "eu", "apac"]),
            wait_policy=MinimumCount(2),
        ),
    ),
    catchup=False,
):
    ...

MinimumCount(-1)는 같은 임계값의 상대적 표현 — "최대 하나 누락" — 이고, 세 멤버 window에 대해 MinimumCount(2)와 동일해요. 기본값에 의존하기보다 의도를 문서화하고 싶을 때 WaitForAll을 명시적으로 전달해요.

세그먼트(범주형) 롤업

버전 3.3.0에 추가됨.

범주형 파티셔닝 — 지역, 테넌트, 실험 변형 — 의 경우 두 프리미티브로 RollupMapper를 구성해요:

  • SegmentWindow(["us", "eu", "apac"])는 하나의 하위 기간을 구성하는 고정 문자열 키 집합을 선언해요. to_upstream은 하위 앵커와 무관하게 전체 집합을 반환해요.
  • FixedKeyMapper("all_regions")는 모든 상위 키를 단일 하위 파티션 키 "all_regions"로 붕괴시켜요.

Scheduler는 upstream producer에서 선언된 모든 세그먼트가 도착할 때까지 하위 Dag run을 보유한 다음 한 번 발화해요. 모든 세그먼트 이벤트가 하나의 AssetPartitionDagRun에 누적되고, 발화된 run의 partition_keyFixedKeyMapper에 전달된 값이에요. 이 구성은 WAIT_FOR_ALL(기본값) 의미 아래에서만 말이 돼요.

from airflow.sdk import (
    DAG,
    Asset,
    FixedKeyMapper,
    PartitionedAtRuntime,
    PartitionedAssetTimetable,
    RollupMapper,
    SegmentWindow,
    asset,
    task,
)

@asset(
    uri="file://incoming/player-stats/multi-region.csv",
    schedule=PartitionedAtRuntime(),
)
def multi_region_player_stats(self, outlet_events):
    # Emit one event per region in a single run.
    outlet_events[self].add_partitions(["us", "eu", "apac"])

# Consumer: fires once all three region partitions have arrived.
with DAG(
    dag_id="segment_region_stats_rollup",
    schedule=PartitionedAssetTimetable(
        assets=Asset.ref(name="multi_region_player_stats"),
        default_partition_mapper=RollupMapper(
            upstream_mapper=FixedKeyMapper("all_regions"),
            window=SegmentWindow(["us", "eu", "apac"]),
        ),
    ),
    catchup=False,
):

    @task
    def aggregate_all_regions(dag_run=None):
        # dag_run.partition_key is the downstream key once all segments arrive.
        print(dag_run.partition_key)

    aggregate_all_regions()

생성은 두 컴포넌트를 모두 검증해요: SegmentWindow는 빈 list, 비문자열 항목, 빈 문자열 키에 대해 ValueError를 발생시키고, 중복 항목은 조용히 제거돼요. FixedKeyMapper는 인자가 비어 있지 않은 문자열이 아니면 ValueError를 발생시켜요. 하나의 consumer Dag가 둘 이상의 asset을 롤업할 때는 뚜렷한 FixedKeyMapper 키를 전달해, 각 롤업이 뚜렷한 버킷을 사용하고 같은 (target_dag_id, partition_key)에서 충돌하지 않게 해요.

런타임에 계산해야 하는 세그먼트 집합이라면 여기에 인코딩하지 마세요 — 대신 consumer측 task에서 완전성을 평가해요 (scheduler는 파티션 집합을 결정하기 위해 사용자 코드를 실행하면 안 되니까요).

런타임에 파티션 키 설정

파티션 키가 미리 알려지지 않을 때 (예: 소스 데이터에서 발견된 watermark, 늦게 도착한 파일, backfill 요청), producer task가 실행 중에 결정하게 해요. Producer를 PartitionedAtRuntime()으로 스케줄링하고 outlet_events[self].add_partitions(...)로 내보낸 이벤트에 키를 기록해요:

from airflow.sdk import PartitionedAtRuntime, asset

@asset(
    uri="file://incoming/player-stats/live-region.csv",
    schedule=PartitionedAtRuntime(),
)
def live_region_player_stats(self, outlet_events):
    # The key is only known once the task runs.
    outlet_events[self].add_partitions("us")

@asset 함수 안에서 self(내보내는 Asset)와 outlet_events(outlet 이벤트 접근자)는 Airflow가 런타임에 채우는 예약 파라미터 이름이에요. 단일 키를 전달하거나, 하나의 run에서 여러 파티션으로 fan out 하려면 리스트를 전달해요. 각 키는 자체 asset 이벤트를 생성하고, 중복 키는 단일 이벤트로 붕괴돼요:

@asset(
    uri="file://incoming/player-stats/multi-region.csv",
    schedule=PartitionedAtRuntime(),
)
def multi_region_player_stats(self, outlet_events):
    outlet_events[self].add_partitions(["us", "eu", "apac"])

런타임 run이 정확히 하나의 파티션 키를 내보내면, producing dag_run.partition_key가 그 키로 back-filled 돼요. 하위 DAG는 timetable 생성 파티션과 같은 방식으로 PartitionedAssetTimetable을 통해 이 이벤트를 소비해요.

Fan-out 매퍼

버전 3.3.0에 추가됨.

FanOutMapperRollupMapper의 거울이에요: 하나의 하위 run이 발화할 수 있을 때까지 여러 상위 이벤트를 보유하는 대신, 단일 상위 이벤트가 window 멤버당 하나의 하위 Dag run으로 fan out 돼요. 상위 key를 window 앵커로 정규화하는 upstream_mapper와, 하위 기간을 열거하는 Window, 그리고 각 window 멤버를 하위 파티션 키 문자열로 변환하는 선택적 downstream_mapper를 구성해요.

시간 windows(WeekWindow, MonthWindow 등)의 경우 기본 downstream_mapper가 자동 적용돼요 — 예를 들어 WeekWindowStartOfDayMapper로 기본 설정되어, 일곱 개의 일별 멤버 각각이 YYYY-MM-DD 문자열로 인코딩돼요. SegmentWindow에는 기본 테이블 항목이 없으므로 명시적 downstream_mapper가 필요해요.

다음 예제는 주간 모델 아티팩트를 일곱 개의 일별 추론 run — 주의 매일 하나의 Dag run — 으로 fan out 해요:

from airflow.sdk import (
    DAG,
    Asset,
    CronPartitionTimetable,
    FanOutMapper,
    PartitionedAssetTimetable,
    StartOfWeekMapper,
    WeekWindow,
    task,
)

weekly_model_artifact = Asset(uri="file://artifacts/models/weekly.bin", name="weekly_model_artifact")

# Producer: emits one partitioned event per week (key is the Monday date).
with DAG(
    dag_id="train_weekly_model",
    schedule=CronPartitionTimetable("0 0 * * 1", timezone="UTC"),
    catchup=False,
):

    @task(outlets=[weekly_model_artifact])
    def train_model():
        pass

    train_model()

# Consumer: one Dag run per day derived from the weekly upstream event.
with DAG(
    dag_id="daily_inference",
    schedule=PartitionedAssetTimetable(
        assets=weekly_model_artifact,
        default_partition_mapper=FanOutMapper(
            upstream_mapper=StartOfWeekMapper(),
            window=WeekWindow(),
            max_downstream_keys=7,
        ),
    ),
    catchup=False,
):

    @task
    def run_inference(dag_run=None):
        # dag_run.partition_key is one daily key, e.g. "2026-03-10".
        print(dag_run.partition_key)

    run_inference()

max_downstream_keys는 하나의 상위 이벤트가 만들 수 있는 하위 Dag run 수를 제한해요. 초과하면 그 이벤트의 run들은 대기열에 넣히지 않고 대신 "partition fan-out exceeded" 감사 로그 항목이 기록돼요. 생략하면 전역 [scheduler] partition_mapper_max_downstream_keys 설정(기본 1000)으로 폴백해요. 의도를 문서화하고 우발적 fan-out 폭발을 막으려면 명시적으로 설정해요.

Window 방향: FORWARD와 BACKWARD

모든 Window는 그 앵커에 상대적으로 window가 어느 기간을 열거할지 제어하는 direction 파라미터를 지원해요.

  • Window.Direction.FORWARD(기본값) — 상위 키에서 시작하는 기간을 결과로 내요. 주간 상위 키 "2026-03-09"(월요일)의 경우 WeekWindow()2026-03-09부터 2026-03-15까지의 일곱 날짜를 결과로 내요.
  • Window.Direction.BACKWARD — 키에서 끝나는 후행 기간을 결과로 내요. 같은 "2026-03-09" 키를 WeekWindow(direction=Window.Direction.BACKWARD)로 하면 그 월요일에 끝나는 일곱 날짜(2026-03-03부터 2026-03-09)를 결과로 내요.
from airflow.sdk import FanOutMapper, PartitionedAssetTimetable, StartOfWeekMapper, WeekWindow, Window

PartitionedAssetTimetable(
    assets=weekly_model_artifact,
    default_partition_mapper=FanOutMapper(
        upstream_mapper=StartOfWeekMapper(),
        window=WeekWindow(direction=Window.Direction.BACKWARD),
    ),
)

방향은 롤업 windows에도 적용돼요 — RollupMapper도 같은 window 클래스를 사용하므로, DayWindow(direction=Window.Direction.BACKWARD)는 하위 키 자정 이전 24시간이 모두 도착할 때까지 하위 run을 보유해요.

플러그인의 커스텀 파티션 매퍼, window, timetable

커스텀 PartitionMapper, Window, 파티션 인지 Timetable 클래스는 각각 AirflowPlugin.partition_mappers, .windows, .timetables에 나열해 Airflow 플러그인으로 제공할 수 있어요. 플러그인이 설치되면 이 클래스들은 core Airflow를 수정하지 않고도 PartitionedAssetTimetableRollupMapper에서 사용할 수 있게 돼요.

커스텀 파티션 매퍼 — 네임스페이스 접두사를 제거해 "eu::daily-sales""us::daily-sales" 같은 상위 키가 하위 키 "daily-sales"로 모두 붕괴되게 해요:

airflow/example_dags/plugins/custom_partition_mapper.py[source]

class PrefixStripMapper(PartitionMapper):
    """
    A partition mapper that strips a fixed namespace prefix from upstream keys.

    Upstream systems often qualify partition keys with a region or environment
    prefix — for example ``"eu::daily-sales"`` or ``"us::daily-sales"``.  A
    downstream asset that aggregates across regions only cares about the base key
    (``"daily-sales"``).  ``PrefixStripMapper`` strips the given prefix (including
    a configurable separator) so that all upstream namespaces collapse to the
    same downstream partition key.

    If the upstream key does not start with the configured prefix the key is
    returned unchanged, which is deliberate: keys that already live in the target
    namespace pass through without modification.

    This class demonstrates registering a custom :class:`PartitionMapper
    <airflow.partition_mappers.base.PartitionMapper>` subclass via the
    ``AirflowPlugin.partition_mappers`` registry. Any plugin that lists it in
    ``partition_mappers = [...]`` makes it available to
    :class:`~airflow.sdk.PartitionedAssetTimetable` and
    :class:`~airflow.partition_mappers.base.RollupMapper` without modifying core
    Airflow.

    :param prefix: The namespace prefix to strip, e.g. ``"eu"``.
    :param separator: The string that separates the prefix from the base key.
        Defaults to ``"::"`` to match a common ``"region::key"`` convention.
    """

    def __init__(
        self,
        prefix: str,
        *,
        separator: str = "::",
        max_downstream_keys: int | None = None,
    ) -> None:
        super().__init__(max_downstream_keys=max_downstream_keys)
        if not prefix:
            raise ValueError("prefix must be a non-empty string.")
        self.prefix = prefix
        self.separator = separator

    def to_downstream(self, key: str) -> str:
        full_prefix = self.prefix + self.separator
        if key.startswith(full_prefix):
            return key[len(full_prefix) :]
        return key

    def serialize(self) -> dict[str, Any]:
        data: dict[str, Any] = {"prefix": self.prefix, "separator": self.separator}
        if self.max_downstream_keys is not None:
            data["max_downstream_keys"] = self.max_downstream_keys
        return data

    @classmethod
    def deserialize(cls, data: dict[str, Any]) -> PartitionMapper:
        return cls(
            prefix=data["prefix"],
            separator=data.get("separator", "::"),
            max_downstream_keys=data.get("max_downstream_keys"),
        )

class PrefixStripMapperPlugin(AirflowPlugin):
    name = "prefix_strip_mapper_plugin"
    partition_mappers = [PrefixStripMapper]

커스텀 롤업 window — 달력 월에서 평일 기간 시작만 결과로 내서 하위 asset이 영업일 상위 파티션만 기다리게 해요:

airflow/example_dags/plugins/business_day_window.py[source]

class BusinessDayWindow(Window):
    """
    A calendar-month rollup window that yields only weekday (Mon–Fri) period-starts.

    The built-in :class:`~airflow.partition_mappers.window.MonthWindow` yields every
    calendar day in the month. ``BusinessDayWindow`` skips Saturdays and Sundays, so
    a monthly downstream asset only waits for the business-day upstream partitions —
    useful for financial or operational pipelines whose upstream data isn't produced
    on weekends.

    This class demonstrates registering a custom ``Window`` subclass via the
    ``AirflowPlugin.windows`` registry: any plugin that lists it in ``windows = [...]``
    makes it available to ``RollupMapper`` without modifying core Airflow.

    *Assumes FORWARD direction and a day-1 month anchor* — the standard contract for
    month-aligned temporal upstream mappers.
    """

    expected_decoded_type: ClassVar[type] = datetime

    def to_upstream(self, period_start: datetime) -> Iterable[datetime]:
        if period_start.day != 1:
            raise ValueError(
                f"BusinessDayWindow expects a period start on day 1 of the month, "
                f"got {period_start.isoformat()}."
            )
        days_in_month = calendar.monthrange(period_start.year, period_start.month)[1]
        return (day for i in range(days_in_month) if (day := period_start + timedelta(days=i)).weekday() < 5)

class BusinessDayWindowPlugin(AirflowPlugin):
    name = "business_day_window_plugin"
    windows = [BusinessDayWindow]

커스텀 런타임 파티션 timetable — partition key를 task 런타임으로 연기하는 스케줄링 가능한 cron timetable. producer task가 파티션을 내보내기 전에 해당 기간의 데이터가 존재하는지 확인할 수 있어요:

airflow/example_dags/plugins/custom_partition_timetable.py[source]

class ScheduledRuntimePartitionTimetable(CronTriggerTimetable):
    """
    A schedulable timetable whose partition key is decided at task runtime.

    Runs fire on the given cron cadence, exactly like an ordinary
    :class:`~airflow.timetables.trigger.CronTriggerTimetable`. The partition key,
    however, is not derived from the schedule: it is set while the producing task
    runs — typically after the task checks whether the period's source data has
    arrived — by calling ``outlet_events[self].add_partitions(...)``.

    This uses runtime partitioning on a regular cron schedule: the timetable stays
    schedulable (``can_be_scheduled`` is ``True``) yet sets
    ``partitioned_at_runtime = True`` so the partition key is deferred to task
    runtime. It differs from :class:`~airflow.sdk.PartitionedAtRuntime` (which also
    defers the key to runtime but never schedules a run on its own) and from
    :class:`~airflow.timetables.trigger.CronPartitionTimetable` (which works out the
    partition key from the cadence ahead of the run). ``partitioned`` stays
    ``False``: no partition key is worked out ahead of the run.

    Registering it via the ``AirflowPlugin.timetables`` registry makes it usable
    by Dag authors without modifying core Airflow.
    """

    partitioned_at_runtime = True

class ScheduledRuntimePartitionTimetablePlugin(AirflowPlugin):
    name = "scheduled_runtime_partition_timetable_plugin"
    timetables = [ScheduledRuntimePartitionTimetable]

Producer Dag는 timetable의 cron 주기로 스케줄링되고 런타임에 파티션 키를 결정해요 — 해당 기간의 상위 데이터가 있을 때만 내보내요. 데이터가 도착하지 않았다면 task는 add_partitions를 호출하지 않으므로(빈·None 키는 어차피 거부됨), 그 기간에 대해 파티션 이벤트 — 따라서 하위 PartitionedAssetTimetable run — 이 생성되지 않아요:

from airflow.sdk import DAG, Asset, task

# ScheduledRuntimePartitionTimetable is provided by the plugin registered above.
from my_plugin.plugins import ScheduledRuntimePartitionTimetable

daily_export = Asset(uri="file://exports/daily.csv", name="daily_export")

with DAG(
    dag_id="export_when_ready",
    schedule=ScheduledRuntimePartitionTimetable("0 6 * * *", timezone="UTC"),
    catchup=False,
):

    @task(outlets=[daily_export])
    def export(*, outlet_events):
        # Fires every day at 06:00 UTC. The partition key is not fixed by the
        # schedule — decide it at runtime from the data itself: find which
        # day's source file has actually landed.
        partition_key = latest_ready_day("s3://raw")  # your own check; e.g. "2026-06-23" or None
        if partition_key:
            build_export(partition_key)
            outlet_events[daily_export].add_partitions(partition_key)
        # If nothing is ready, emit nothing: with no partition key recorded,
        # no partitioned event is produced and no downstream
        # PartitionedAssetTimetable run is triggered for this period.

    export()

완전한 실행 가능 예제는 airflow-core/src/airflow/example_dags/example_asset_partition.py를 참고해요.

더 알아보기 (Learn more)