DAG Result
DAG Result
이 페이지는 버전 3.3에 추가된 DAG Result(실험적 기능)를 다뤄요. DAG의 하나 이상의 Task를 result task로 지정하면, 그 반환값을 /dags/{dag_id}/dagRuns/{dag_run_id}/wait API가 노출해서, 폴링이나 커스텀 글루 코드 없이 Airflow DAG를 API 엔드포인트·채팅 에이전트·추론 서비스에 내장하기 쉽게 해줘요.
출처: 문서
본문
이것은 실험적 기능이에요.
버전 3.3에 추가됨.
Airflow는 DAG의 하나 이상의 Task를 result task로 지정하는 것을 지원해요. DAG에 result task가 있으면 그 반환값이 /dags/{dag_id}/dagRuns/{dag_run_id}/wait API에 의해 노출되어, 폴링이나 커스텀 글루 코드 없이 응답값이 필요한 API 엔드포인트, 채팅 에이전트, 추론 서비스에 Airflow DAG를 손쉽게 내장할 수 있어요.
결과 Task 표시하기
result task를 지정하는 동등한 두 가지 방법이 있어요.
@result 사용하기 (TaskFlow)
airflow.sdk.result() 데코레이터는 TaskFlow Task를 DAG의 result로 표시해요. 반드시 @task 데코레이터 위에 적용해야 해요:
from airflow.sdk import dag, result, task
@dag
def my_dag():
@task
def fetch_data():
return {"answer": 42}
@result
@task
def compute(data):
return data["answer"] * 2
compute(fetch_data())
my_dag()
@result로 데코레이트된 Task는 매핑된 Task로 expand()나 expand_kwargs()할 수 있어요.
@dag 함수에서 반환하기
@dag 데코레이터를 사용할 때 함수 본문에서 Task의 XComArg를 직접 반환하면 그 Task가 자동으로 result로 지정돼요. 명시적 @result가 필요 없어요:
from airflow.sdk import dag, task
@dag
def my_dag():
@task
def fetch_data():
return {"answer": 42}
@task
def compute(data):
return data["answer"] * 2
return compute(fetch_data()) # 'compute' marked as the result task.
my_dag()
일반 XComArg(@task 함수의 직접 호출 결과)만 허용돼요. Airflow는 일반 Python 객체 같은 다른 반환값은 조용히 무시하므로, 리터럴 정수나 문자열을 반환해도 효과가 없어요.
Note
일반
return_valueXCom(즉 Python 함수가 반환하는 것)만 DAG result로 지정될 수 있어요. 다른 키로 밀린 XCom이나multiple_outputs=True로 생성된 XCom은 자격이 없어요.
결과 가져오기
DAG에 result task가 있으면 GET /api/v2/dags/{dag_id}/dagRuns/{dag_run_id}/wait 엔드포인트를 호출해 run이 끝날 때까지 블로킹하고 단일 요청으로 결과를 수집해요:
curl -s "http://localhost:8080/api/v2/dags/my_dag/dagRuns/<run_id>/wait"
이 엔드포인트는 새 줄로 구분된 JSON(NDJSON)을 스트리밍해요. 각 줄은 현재 DAG run 상태를 보고하고, 마지막 줄은 result task의 반환값을 task ID별로 담은 results 키를 포함해요:
{"state": "running"}
{"state": "success", "results": {"compute": 84}}
result 쿼리 파라미터는 기본 동작을 오버라이드하고, result task가 선언됐는지와 무관하게 호출자가 어떤 task ID에서 XCom을 수집할지 선택할 수 있게 해줘요:
# Collect XCom from a specific task regardless of @result marking
curl -s "http://localhost:8080/api/v2/dags/my_dag/dagRuns/<run_id>/wait?result=some_task_id"
# Suppress all result collection
curl -s "http://localhost:8080/api/v2/dags/my_dag/dagRuns/<run_id>/wait?result="
result를 생략하면 API는 작성자가 선언한 result task의 XCom을 반환해요. result를 하나 이상의 task ID로 설정하면 그 task ID들이 대신 사용돼요. result를 빈 문자열로 설정하면 XCom이 수집되지 않아요.
result task가 동적으로 매핑되면 results의 항목은 task ID와 map index 순으로 정렬된, 매핑된 인스턴스당 하나의 값이 담긴 리스트예요.