매일 원본 데이터를 읽고, 가공한 뒤, 결과를 저장하는 작업을 스크립트 하나로 시작할 수 있다. 하지만 실패한 단계만 다시 실행하거나 과거 기간을 재처리하고, 여러 작업의 의존성을 추적하려면 실행 관리 도구가 필요해진다.
Airflow는 이런 유한한 워크플로우의 실행 순서와 상태를 관리하는 오케스트레이터다. 데이터를 직접 처리하는 엔진과는 역할이 다르다. SQL, Python, Spark 작업 등을 태스크로 연결할 수 있지만 실제 처리 방식과 데이터 저장은 각 태스크가 담당한다.
이 글은 Airflow 3.0.x와 Task SDK를 기준으로 한다. Airflow 2의 import 경로나 schedule_interval 예제를 그대로 섞지 않는다.
DAG는 태스크와 의존성을 정의한 방향성 비순환 그래프다. 순환이 없다는 것은 같은 실행 안에서 의존성이 끝없이 돌아오지 않는다는 뜻이다. 같은 DAG를 매일 실행하는 것은 가능하다.
| 개념 | 의미 |
|---|---|
| DAG | 작업 구조와 실행 규칙의 정의 |
| DAG Run | 특정 시점·조건에서 생성된 DAG의 한 번의 실행 |
| Task | DAG에 정의된 작업 단위 |
| Task Instance | 특정 DAG Run에서의 태스크 실행 상태 |
| Operator·TaskFlow task | 작업 동작을 정의하는 방법 |
extract >> transform은 실행 의존성을 연결한다. 기본 all_success 규칙에서는 상위 태스크가 성공해야 하위 태스크를 실행하지만, 다른 trigger rule을 사용하면 조건이 달라진다. 화살표가 데이터를 자동으로 전달한다는 뜻도 아니다.
TaskFlow API에서는 앞 태스크의 반환값을 다음 태스크의 인자로 연결할 수 있다. 이 연결에는 의존성과 XCom이 사용된다. XCom에는 작은 메타데이터를 전달하고, 큰 데이터셋은 외부 저장소에 둔 뒤 URI나 작업 ID를 전달하는 편이 적합하다.
다음 파일은 Airflow 3 환경의 DAG 폴더에 넣는 예제다. 입력을 고정해 외부 API나 별도 /scripts 파일 없이 흐름을 확인할 수 있게 했다. 수동 실행에서 합계 250을 로그로 남기며 실제 DB 적재는 하지 않는다.
# sales_demo.py
from datetime import timedelta
import pendulum
from airflow.sdk import dag, task
@dag(
dag_id="sales_demo",
schedule=None,
start_date=pendulum.datetime(2025, 1, 1, tz="UTC"),
catchup=False,
tags=["example"],
)
def sales_demo():
@task(retries=2, retry_delay=timedelta(seconds=10))
def extract() -> list[dict]:
return [
{"order_id": "A", "amount": 100},
{"order_id": "B", "amount": 150},
]
@task
def transform(rows: list[dict]) -> int:
if any(type(row["amount"]) is not int or row["amount"] < 0 for row in rows):
raise ValueError("amount must be a non-negative integer")
return sum(row["amount"] for row in rows)
@task
def load(total: int) -> None:
print(f"total={total}")
load(transform(extract()))
sales_dag = sales_demo()
Airflow가 구성된 Linux 환경에서 DAG를 등록하고 UI 또는 CLI로 수동 실행한다. Windows에서는 지원되는 WSL2·Linux 컨테이너 등의 환경을 준비한다. 일반 Python으로 파일을 실행하는 것과 Airflow 런타임에서 DAG Run을 실행하는 것은 다르다.
airflow dags list
airflow dags trigger sales_demo
retries=2는 최초 실행 뒤 최대 두 번 더 시도한다는 뜻이다. 입력 자체가 잘못된 경우와 일시적인 연결 오류를 같은 재시도 정책으로 처리할지는 별도로 판단한다.
| 구성 요소 | 주요 역할 |
|---|---|
| DAG Processor | DAG 파일을 읽고 정의를 처리 |
| Scheduler | 실행 가능 조건을 판단하고 태스크 실행을 조정 |
| Executor·Worker | 선택한 방식으로 태스크 실행 |
| API Server·Task Execution API | UI·클라이언트 및 태스크 런타임의 API 통신 |
| Metadata DB | DAG Run·Task Instance 등 제어 상태 저장 |
| Triggerer | deferrable 작업의 비동기 대기 지원 |
Airflow 3에서는 태스크 코드가 Airflow 메타데이터 DB 세션이나 내부 모델에 직접 접근하는 방식이 제한된다. Task SDK와 Task Execution API를 통해 실행 상태·XCom 등 런타임 자원에 접근한다. 업무용 DB에 연결해 데이터를 처리하는 것과 Airflow 내부 DB에 접근하는 것은 구분해야 한다.
DAG를 해석하는 시점과 태스크가 실행되는 시점도 다르다. 파일 최상위에서 대용량 데이터 조회나 외부 API 호출을 하면 DAG 처리 과정에 부담을 준다. 실제 처리는 태스크 함수 안에 두고, DAG 정의 코드는 가볍게 유지한다.
위 예제의 schedule=None은 자동 주기 실행을 만들지 않는다. 매일 처리하는 파이프라인으로 확장하려면 스케줄과 함께 어느 시간 구간의 데이터를 처리하는지 정한다.
시간 기반 스케줄에서는 실행 시각과 논리적인 데이터 구간이 다를 수 있다. 자정에 시작한 실행이 막 끝난 하루를 처리할 수도 있다. 태스크에서 무조건 현재 시각을 기준으로 조회하면 과거 실행을 다시 돌릴 때 다른 데이터를 읽게 된다.
업무 날짜·시작과 종료 시각을 명시하고, 같은 구간을 재실행해도 같은 입력을 읽게 한다. 표준 시간대와 일광 절약 시간제 영향을 함께 고려한다. start_date를 매번 현재 시각으로 바꾸는 방식도 피한다.
catchup은 스케줄에 따라 과거 구간의 실행을 만들지 결정하는 설정이다. 운영자가 특정 기간을 다시 처리하는 backfill과 목적을 구분한다. 과거 데이터가 수정됐는지, 이미 처리한 구간을 어떻게 덮어쓸지까지 정해야 재처리가 의미를 갖는다.
외부 DB에 결과를 저장한 뒤 태스크 완료 기록을 남기기 전에 프로세스가 종료될 수 있다. 재실행하면 같은 작업이 또 수행된다. 실행 도구가 재시도를 제공한다는 사실만으로 외부 업무 효과가 한 번만 발생하지는 않는다.
예를 들어 (업무 날짜, 집계 종류)를 결과의 유일한 키로 정해 upsert하거나, 동일 입력을 처리한 기록을 확인하는 방식이 필요하다. 이메일·결제 같은 부수 효과는 별도의 멱등 식별자를 둔다. 실패를 복구하려고 태스크 상태를 지우기 전에 이미 외부에 반영된 결과가 있는지 확인한다.
병렬 처리를 늘릴 때도 대상 DB의 연결 수와 API 제한을 고려한다. Pool과 동시 실행 제한은 오케스트레이터뿐 아니라 외부 시스템을 보호하는 기준으로 설정한다.
정기 배치, 여러 시스템의 작업 연결, 기간별 재처리와 운영 이력 추적이 중요하다면 Airflow를 검토할 만하다. 다양한 provider가 있지만 필요한 provider 패키지와 버전은 별도로 관리한다.
상시 실행되는 무한 스트림 처리나 매우 낮은 지연의 개별 이벤트 처리는 별도의 처리 엔진이 더 적합할 수 있다. Airflow에서 그런 엔진의 작업을 시작·감시하는 것과 Airflow 자체를 스트림 엔진으로 사용하는 것은 다르다.
Airflow에도 동적 태스크 매핑과 데이터 품질 검증을 태스크로 구성하는 방법이 있다. 따라서 “동적인 작업이나 품질 검사는 불가능하다”고 단정하지 않는다. 대신 팀의 배치 운영 경험, 실행 빈도, 복잡도와 인프라 비용에 맞춰 판단한다.
운영에서는 DAG 등록 성공만 보지 말고 실행 지연, 태스크 실패·재시도, 큐 대기, 오래 실행 중인 작업과 외부 자원 부하를 살핀다. DAG 정의, 실제 실행, 외부 결과의 정확성을 각각 확인하는 것이 안정적인 재처리의 출발점이다.