에셋 기반 스케줄링
에셋 기반 스케줄링 (Asset-Aware Scheduling)
배치 워크플로가 꼭 시간에 맞춰 돌아야 하는 것은 아니에요. "어떤 파일이 새로 생겼을 때", "어떤 태스크가 데이터를 갱신했을 때" 이어서 다음 파이프라인이 자동으로 돌게 하고 싶을 때가 있죠. Airflow의 **에셋(Asset)**을 쓰면 태스크가 데이터를 갱신하는 시점에 다음 DAG을 트리거할 수 있어요. 시간 스케줄과 달리 데이터 상태에 반응하는 스케줄링이라, 데이터 파이프라인 사이의 의존성을 자연스럽게 표현해 줘요. (버전 2.4부터 추가됐어요.)
본문
가장 간단한 형태는 태스크가 outlets로 에셋을 갱신하고, 다른 DAG이 그 에셋을 schedule로 기다리는 거예요.
from airflow.sdk import DAG, Asset
with DAG(...):
MyOperator(
# this task updates example.csv
outlets=[Asset("s3://asset-bucket/example.csv")],
...,
)
with DAG(
# this Dag should be run when example.csv is updated (by dag1)
schedule=[Asset("s3://asset-bucket/example.csv")],
...,
):
...
에셋으로 DAG 스케줄하기
에셋은 DAG의 데이터 의존성을 명시할 때 써요. 아래 예시는 producer DAG의 producer 태스크가 성공적으로 완료되면 consumer DAG을 스케줄해요. Airflow는 태스크가 성공했을 때만 에셋을 updated로 표시해요. 태스크가 실패하거나 스킵되면 갱신이 일어나지 않고 consumer DAG도 스케줄되지 않아요.
example_asset = Asset("s3://asset/example.csv")
with DAG(dag_id="producer", ...):
BashOperator(task_id="producer", outlets=[example_asset], ...)
with DAG(dag_id="consumer", schedule=[example_asset], ...):
...
에셋과 DAG 사이의 관계 목록은 Asset Views에서 볼 수 있어요.
여러 에셋
schedule 파라미터가 리스트이므로 DAG은 여러 에셋을 요구할 수 있어요. Airflow는 DAG이 소비하는 모든 에셋이, 마지막으로 DAG이 실행된 이후에 최소 한 번씩 갱신된 뒤에 DAG을 스케줄해요.
with DAG(
dag_id="multiple_assets_example",
schedule=[
example_asset_1,
example_asset_2,
example_asset_3,
],
...,
):
...
트리거한 에셋 이벤트에서 정보 가져오기
트리거된 DAG은 triggering_asset_events 템플릿(또는 파라미터)으로 자신을 트리거한 에셋의 정보를 가져올 수 있어요. 이 값은 이런 모양의 딕셔너리예요. 에셋을 키로, 그 에셋이 발생시킨 이벤트 리스트를 값으로 가집니다.
{
Asset("s3://asset-bucket/example.csv"): [
AssetEvent(uri="s3://asset-bucket/example.csv", source_dag_run=DagRun(...), ...),
...,
],
Asset("s3://another-bucket/another.csv"): [
AssetEvent(uri="s3://another-bucket/another.csv", source_dag_run=DagRun(...), ...),
...,
],
}
Jinja로 접근하기 — 주의할 점이 하나 있어요. Airflow 3.1.x에서는 AssetEvent의 source_dag_run 속성이 Jinja 템플릿에 노출되지 않아서 접근하면 템플릿 렌더링 에러가 나요. (3.2.0부터 추가됨.) 트리거링 에셋 이벤트의 정보를 오퍼레이터로 전달하려면 Jinja 템플릿을 쓸 수 있어요.
단일 에셋으로 트리거된 경우, 예를 들어 Snowflake 테이블을 쿼리하는 DAG에서 갱신된 데이터 구간만 읽고 싶다면 이렇게 접근해요. triggering_asset_events.values() | first | first는 ① 모든 에셋 이벤트 리스트를 얻고(값들) ② 첫 번째 리스트를 얻고(단일 트리거 에셋이므로) ③ 그 리스트에서 첫 이벤트를 얻는 식으로 동작해요.
example_snowflake_asset = Asset("snowflake://my_db/my_schema/my_table")
with DAG(dag_id="query_snowflake_data", schedule=[example_snowflake_asset], ...):
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 }}';
""",
)
여러 에셋으로 트리거되는 경우에는 Jinja 템플릿에서 반복문으로 돌리면 돼요.
with DAG(dag_id="process_assets", schedule=[asset1, asset2], ...):
BashOperator(
task_id="process",
bash_command="""
{% for asset_uri, events in triggering_asset_events.items() %}
echo "Processing asset: {{ asset_uri }}"
{% for event in events %}
echo " Triggered by DAG: {{ event.source_dag_run.dag_id }}"
echo " Data interval start: {{ event.source_dag_run.data_interval_start }}"
echo " Data interval end: {{ event.source_dag_run.data_interval_end }}"
{% endfor %}
{% endfor %}
""",
)
Python으로 접근하기 — TaskFlow 태스크에서는 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():
print(f"Asset: {asset.uri}")
for event in asset_events:
print(f" - Triggered by DAG run: {event.source_dag_run.dag_id}")
print(
f" Data interval: {event.source_dag_run.data_interval_start} to {event.source_dag_run.data_interval_end}"
)
print(f" Run ID: {event.source_dag_run.run_id}")
print(f" Timestamp: {event.timestamp}")
print_triggering_asset_events()
이벤트 기반 스케줄링
푸시 기반 (REST API) — 외부 시스템이 REST API로 Airflow에 에셋 이벤트를 밀어 넣을(push) 수 있어요. 예를 들어 DAG waiting_for_asset_1_and_2는 "asset-1"과 "asset-2" 둘 모두가 갱신된 태스크에 의해 트리거돼요. "asset-1"이 갱신되면 Airflow는 레코드를 하나 만들어 두고, "asset-2"가 갱신될 때 DAG을 트리거할 수 있게 해 둬요. 이런 레코드를 queued asset events라고 불러요.
with DAG(
dag_id="waiting_for_asset_1_and_2",
schedule=[Asset("asset-1"), Asset("asset-2")],
...,
):
...
queuedEvent API 엔드포인트로 이 레코드들을 다룰 수 있어요.
- DAG용 queued asset event 조회:
/assets/queuedEvent/{uri} - DAG용 queued asset event 목록 조회:
/dags/{dag_id}/assets/queuedEvent - DAG용 queued asset event 삭제:
/assets/queuedEvent/{uri} - DAG용 queued asset event 목록 삭제:
/dags/{dag_id}/assets/queuedEvent - 에셋용 queued asset event 목록 조회:
/dags/{dag_id}/assets/queuedEvent/{uri} - 에셋용 queued asset event 삭제:
DELETE /dags/{dag_id}/assets/queuedEvent/{uri}
자세한 사용법과 파라미터는 Airflow API 문서를 확인하세요.
풀 기반 (Asset Watchers) — 외부 시스템이 이벤트를 push하는 대신 Airflow가 외부 이벤트 소스에서 직접 끌어올(pull) 수도 있어요. AssetWatcher 클래스와 이벤트 기반 스케줄링 호환 트리거로 구현해요. AssetWatcher는 큐나 스토리지 같은 외부 소스를 모니터링하고, 관련 이벤트가 발생하면 해당 에셋을 갱신해 DAG 실행을 트리거해요. 무한 재스케줄링을 막기 위해 BaseEventTrigger를 상속한 트리거만 호환돼요.
조건부 표현식으로 고급 에셋 스케줄링
Airflow는 에셋과 함께 조건부 표현식을 쓰는 고급 스케줄링 기능도 제공해요. 논리 연산자로 에셋 갱신에 기반한 복잡한 DAG 실행 의존성을 정의할 수 있죠. 지원하는 연산자는 두 가지예요.
- AND (
&): 지정된 모든 에셋이 갱신된 후에만 DAG을 트리거. - OR (
|): 지정된 에셋 중 하나라도 갱신되면 DAG을 트리거.
이 연산자들로 워크플로가 더 동적이고 유연한 에셋 갱신 조건을 갖출 수 있어요.
두 에셋이 모두 갱신된 후에만 실행하려면 AND를 써요.
dag1_asset = Asset("s3://dag1/output_1.txt")
dag2_asset = Asset("s3://dag2/output_1.txt")
with DAG(
# Consume asset 1 and 2 with asset expressions
schedule=(dag1_asset & dag2_asset),
...,
):
...
어느 하나라도 갱신되면 실행하려면 OR를 써요.
with DAG(
# Consume asset 1 or 2 with asset expressions
schedule=(dag1_asset | dag2_asset),
...,
):
...
괄호로 우선순위를 조합할 수도 있어요. 예를 들어 "에셋 1, 또는 (에셋 2와 3 모두)" 같은 조건이요.
dag3_asset = Asset("s3://dag3/output_3.txt")
with DAG(
# Consume asset 1 or both 2 and 3 with asset expressions
schedule=(dag1_asset | (dag2_asset & dag3_asset)),
...,
):
...
에셋 앨리어스 기반 스케줄링
작성 문법에서 AssetAlias를 이름으로 참조하면 그와 연결된 에셋 이벤트를 스케줄링에 사용해요. DAG은 outlets=AssetAlias("xxx") 태스크가, 앨리어스가 Asset("s3://bucket/my-task")로 해석될 때에만 트리거될 수 있어요. 런타임에 outlet AssetAlias("out")인 태스크가 에셋 하나 이상과 연결되면 DAG은 그 에셋의 정체와 무관하게 실행돼요. 특정 태스크 실행에서 앨리어스에 연결된 에셋이 없다면 하위 DAG은 트리거되지 않아요.
with DAG(dag_id="asset-producer"):
@task(outlets=[Asset("example-alias")])
def produce_asset_events():
pass
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"))
with DAG(dag_id="asset-consumer", schedule=Asset("s3://bucket/my-task")):
...
with DAG(dag_id="asset-alias-consumer", schedule=AssetAlias("example-alias")):
...
에셋 앨리어스는 DAG 파싱 중에 에셋으로 해석돼요. 그래서 min_file_process_interval 설정이 크면 앨리어스가 해석되지 않을 가능성이 있어요. 그 경우 DAG 파싱을 수동으로 트리거해서 해결할 수 있어요.
위 예시에서 asset-alias-producer가 실행되면 앨리어스 AssetAlias("example-alias")는 Asset("s3://bucket/my-task")로 해석돼요. 그런데 asset-alias-consumer DAG은 다음 재파싱까지 기다려야 스케줄이 갱신돼요. 이 문제를 위해 Airflow는 앨리어스가 이전에 의존하지 않던 에셋으로 해석되면 그 앨리어스에 의존하는 DAG을 재파싱해요. 그래서 asset-producer 실행 후 asset-consumer와 asset-alias-consumer 두 DAG이 모두 트리거돼요.
에셋과 시간 스케줄 결합
에셋 이벤트와 시간 기반 스케줄을 동시에 쓰려면 AssetOrTimeSchedule로 DAG을 예약할 수 있어요. 데이터 갱신에도 트리거되고 고정 timetable대로 주기적으로도 실행돼야 하는 워크플로에 적합해요. 자세한 내용은 Timetable 페이지의 AssetOrTimeSchedule 섹션을 참고하세요.