Asset 기반 스케줄링 기초 (schedule=[asset]과 트리거 이벤트 조회)
Apache Airflow Official Documentation (in-repo snapshot) — Apache Software Foundation authoring-and-scheduling/asset-scheduling.rst - Quickstart, Schedule Dags with assets, Multiple assets, Fetching information from a triggering asset event (versionadded 2.4)
이 모듈을 다 읽으면
- schedule 파라미터에 Asset을 지정해 Dag를 Asset 갱신 기반으로 트리거하는 방법을 설명할 수 있다
- 태스크 실패·스킵 시 Asset이 갱신되지 않는다는 규칙을 설명할 수 있다
- 여러 Asset을 소비하는 Dag가 정확히 언제 트리거되는지 판단할 수 있다
- triggering_asset_events로 트리거를 유발한 Asset 정보를 Jinja와 Python에서 각각 조회하는 방법을 설명할 수 있다
Dag의 schedule 인자에 Asset(들)을 지정하면 시간이 아니라 데이터 갱신을 기준으로 Dag를 스케줄링할 수 있다. 태스크가 성공적으로 완료된 경우에만 Asset이 갱신된 것으로 표시되며, 여러 Asset을 소비하는 Dag는 마지막 실행 이후 모든 Asset이 최소 한 번씩 갱신되어야 트리거된다. triggering_asset_events를 통해 어떤 이벤트가 이번 실행을 유발했는지 Jinja 템플릿이나 Python 함수에서 조회할 수 있다.
Asset 기반 스케줄링 퀵스타트
producer 태스크가 outlets로 Asset을 갱신하고, consumer Dag는 ``schedule``에 그 Asset(들)을 지정한다.
from airflow.sdk import DAG, Asset
with DAG(...):
MyOperator(
outlets=[Asset("s3://asset-bucket/example.csv")],
...,
)
with DAG(
schedule=[Asset("s3://asset-bucket/example.csv")],
...,
):
...
Airflow는 태스크가 '성공적으로 완료된 경우에만' Asset을 ``updated``로 표시한다. 태스크가 실패하거나 스킵되면 갱신이 일어나지 않고, 그 Asset을 소비하는 Dag도 스케줄되지 않는다. Asset과 Dag 사이의 이런 관계는 Asset Views 화면에서 목록으로 확인할 수 있다.
핵심 포인트
- consumer Dag의 schedule에 Asset을 지정하면 시간이 아니라 데이터 갱신을 기준으로 스케줄링되며, producer 태스크는 outlets로 Asset을 갱신한다
- 태스크가 성공적으로 완료된 경우에만 Asset이 updated로 표시되고, 실패·스킵 시에는 갱신이 일어나지 않아 소비 Dag도 트리거되지 않는다
여러 Asset 소비하기
``schedule``은 리스트이므로 여러 Asset을 요구하는 Dag를 만들 수 있다. Airflow는 Dag가 소비하는 '모든' Asset이 마지막 실행 이후 최소 한 번씩 갱신되어야 다음 Dag run을 스케줄한다.
with DAG(
dag_id="multiple_assets_example",
schedule=[example_asset_1, example_asset_2, example_asset_3],
...,
):
...
한 Asset이 다른 Asset들보다 먼저, 그리고 여러 번 갱신되더라도 모든 소비 Asset이 갱신되기 전까지 다운스트림 Dag run은 만들어지지 않으며, 조건이 충족되는 순간 단 한 번만 실행된다. 즉 갱신 횟수가 아니라 '마지막 실행 이후 전부 최소 1회 갱신됐는가'가 트리거 조건이다.
핵심 포인트
- schedule에 여러 Asset을 리스트로 지정하면, 마지막 실행 이후 그 Asset들이 모두 최소 한 번씩 갱신되어야 다운스트림 Dag가 트리거된다
- 한 Asset이 다른 Asset보다 먼저 여러 번 갱신되더라도 다운스트림은 조건이 충족되는 시점에 딱 한 번만 실행된다
트리거한 Asset 이벤트 정보 조회하기
트리거된 Dag는 ``triggering_asset_events`` 템플릿/파라미터로 자신을 트리거한 Asset의 정보를 읽을 수 있다. 이는 ``{Asset: [AssetEvent, ...], ...}`` 형태의 딕셔너리다.
Jinja에서 접근하는 기능은 3.2.0에서 추가되었다(3.1.x에서는 ``AssetEvent``의 ``source_dag_run`` 속성이 Jinja 템플릿에 노출되지 않아 접근 시 렌더링 오류가 난다).
SQLExecuteQueryOperator(
task_id="query",
conn_id="snowflake_default",
sql='''
SELECT * FROM my_db.my_schema.my_table
WHERE "updated_at" >= '{{ (triggering_asset_events.values() | first | first).source_dag_run.data_interval_start }}'
AND "updated_at" < '{{ (triggering_asset_events.values() | first | first).source_dag_run.data_interval_end }}';
''',
)
``triggering_asset_events.values() | first | first``는 (1) 모든 Asset의 이벤트 리스트들을 가져오고, (2) 그중 첫 리스트를 가져오고(트리거 Asset이 하나뿐이라 가정), (3) 그 리스트의 첫 ``AssetEvent``를 가져오는 순서로 동작한다. 여러 Asset이 트리거한 경우에는 ``{% for asset_uri, events in triggering_asset_events.items() %}``로 순회할 수 있다.
Python 함수에서는 ``triggering_asset_events``를 파라미터로 받아 직접 순회한다.
@task
def print_triggering_asset_events(triggering_asset_events=None):
if triggering_asset_events:
for asset, asset_events in triggering_asset_events.items():
for event in asset_events:
print(event.source_dag_run.dag_id, event.source_dag_run.data_interval_start, event.timestamp)
하나의 Asset에 대해서도 여러 이벤트가 있을 수 있으므로(예: 여러 번 갱신된 뒤에 소비 Dag가 실행되는 경우), 그 이벤트들을 어떻게 처리할지(마지막 실행 이후 신규 데이터를 전부 처리할지, 트리거 이벤트 각각을 개별 처리할지)는 Dag 작성자가 결정해야 한다.
핵심 포인트
- triggering_asset_events는 {Asset: [AssetEvent, ...]} 형태의 딕셔너리이며, Jinja에서 AssetEvent.source_dag_run에 접근하는 기능은 3.2.0에 추가되었다(3.1.x는 미지원)
- Python 태스크에서는 triggering_asset_events를 파라미터로 받아 dict를 직접 순회하며, 한 Asset에 여러 이벤트가 있을 때 이를 어떻게 처리할지는 Dag 작성자의 몫이다