사전 Apache Airflow
구현체

Apache Airflow

gabury1고친 사람 github-actions[bot]

Apache Airflow 는 여러 단계로 된 일을 정해진 때에 정해진 순서대로 돌려 주는 도구입니다. 어느 단계가 어느 단계 뒤에 와야 하는지를 코드로 적어 두면 그 순서를 지키며 실행합니다. 실패한 단계는 다시 돌립니다. 무엇이 언제 성공했고 실패했는지도 화면에 남깁니다.

쉽고 빠른 이해

여러 단계로 된 정기적인 일을 순서대로 돌려 주는 도구입니다. 밤마다 「어제 주문을 꺼낸다 → 합계를 낸다 → 보고서 표에 넣는다」를 돌리는 일이 그런 일입니다.

왜 이렇게 하나. 단계마다 시각만 맞춰 따로 걸어 두면 앞 단계가 늦거나 실패해도 뒤 단계가 그냥 돕니다. 어디서 멈췄는지도 로그를 뒤져야 압니다. Airflow 는 단계 사이의 앞뒤를 알고 있어서 앞 단계가 성공해야 뒤 단계를 돌립니다.

어떻게 도나.

  1. 단계와 그 순서를 파이썬 파일 하나에 적습니다.
  2. 스케줄러(때가 된 단계를 골라 내는 프로그램)가 그 파일을 읽고 때가 된 단계를 실행하도록 넘깁니다.
  3. 단계마다 성공과 실패를 데이터베이스에 적습니다. 실패하면 정한 횟수만큼 다시 돌립니다.

대가. 스케줄러와 데이터베이스, 웹 화면을 따로 띄워 굴려야 합니다. 단계 하나를 띄우는 데에도 절차가 있어서 몇 초 안에 끝나야 하는 일에는 부담이 됩니다. 끝없이 흘러드는 데이터를 바로바로 처리하는 일도 맡지 않습니다.

상세

이 절은 Airflow 가 맡는 일과 cron 으로 모자란 대목을 본 뒤 파이프라인을 적는 코드와 그 코드를 돌리는 프로그램을 따라갑니다.

Airflow 는 워크플로를 정해진 때에 정해진 순서로 돌려 주는 오픈 소스 도구입니다. 워크플로는 순서가 정해진 작업 여러 개의 묶음입니다. 밤마다 주문 데이터베이스에서 어제 주문을 꺼내 합계를 내고 보고서 표에 넣는 일이 그런 묶음입니다. 데이터 쪽에서는 이런 묶음을 파이프라인이라고 부릅니다.

Airflow 가 흔히 맡는 일은 ETL(Extract-Transform-Load, 추출-변환-적재)입니다. 여러 원천에서 데이터를 꺼내 모양을 바꾼 뒤 한곳에 싣는 일입니다. 하루에 한 번, 한 시간에 한 번처럼 주기를 두고 도는 배치 처리가 대부분입니다.

이름 앞의 Apache 는 이 프로젝트를 아파치 소프트웨어 재단이 관리한다는 표시입니다. 흔히 줄여서 Airflow 라고 부릅니다.

cron 으로 모자란 대목

cron 은 유닉스 계열 운영체제에서 명령을 정해진 시각에 돌려 주는 프로그램입니다. 「매일 새벽 두 시에 이 스크립트를 돌려라」 같은 줄을 적어 두면 그 시각에 실행합니다. 작업이 하나뿐이면 이것으로 충분합니다.

작업이 셋으로 늘어 앞뒤가 생기면 사정이 달라집니다. cron 은 작업끼리의 순서를 모릅니다. 그래서 추출은 두 시, 변환은 세 시, 적재는 네 시처럼 시각을 벌려 걸게 됩니다. 추출이 한 시간을 넘기는 날에는 변환이 덜 채워진 데이터를 읽습니다.

실패도 따로 챙겨야 합니다. cron 은 실패한 작업을 다시 돌리지 않습니다. 어느 날 어느 단계가 실패했는지 보려면 서버에 들어가 로그를 뒤져야 합니다.

Airflow 는 이 세 가지를 도구 안에 넣었습니다. 작업 사이의 앞뒤를 알고 있으므로 앞 작업이 성공해야 뒤 작업을 시작합니다. 실패한 작업은 정한 횟수만큼 재시도합니다. 날마다 어느 작업이 어떻게 끝났는지는 데이터베이스에 남고 웹 화면에서 봅니다.

DAG 로 적는 파이프라인

Airflow 에서 파이프라인 하나는 DAG(Directed Acyclic Graph, 방향 비순환 그래프) 하나입니다. DAG 는 점을 화살표로 이은 그림입니다. 화살표를 아무리 따라가도 출발한 점으로 돌아오지 않습니다. 그래서 「A 가 끝나야 B, B 가 끝나야 A」 같은 꼬인 순서가 생기지 않습니다.

DAG 의 점 하나가 태스크입니다. 앞에서 단계·작업이라 부른 것을 Airflow 는 태스크라 부릅니다. 셸 명령 하나, 파이썬 함수 하나, SQL(Structured Query Language) 질의 하나가 각각 태스크가 됩니다. 점을 잇는 화살표는 「이 태스크 다음에 저 태스크」라는 순서입니다.

태스크는 오퍼레이터로 만듭니다. 오퍼레이터는 자주 하는 일을 미리 만들어 둔 태스크 틀입니다. 셸 명령을 돌리는 오퍼레이터에 명령만 넣으면 태스크 하나가 됩니다. 틀이 있으니 실행과 로그 남기기를 매번 새로 짜지 않아도 됩니다.

DAG 는 파이썬 파일로 적습니다. 아래는 주문과 환불을 따로 꺼낸 뒤 둘을 합쳐 보고서를 만드는 DAG 입니다. 가져오기 줄은 뺐습니다.

Python
with DAG(
    dag_id="daily_sales",
    start_date=datetime(2026, 9, 1),
    schedule="@daily",  # 하루 한 번
    catchup=False,      # 지난날은 안 채움
):
    orders = BashOperator(
        task_id="extract_orders",
        bash_command="python orders.py",
    )
    refunds = BashOperator(
        task_id="extract_refunds",
        bash_command="python refunds.py",
    )
    report = BashOperator(
        task_id="build_report",
        bash_command="python report.py",
    )
    [orders, refunds] >> report  # 둘 다음에

위쪽 with DAG(...) 는 이 파이프라인의 이름과 도는 주기를 정합니다. 아래쪽에서는 BashOperator 로 태스크 셋을 만듭니다. BashOperator 는 셸 명령을 돌리는 오퍼레이터입니다.

마지막 줄의 >> 가 순서를 적습니다. 왼쪽이 끝나야 오른쪽이 돕니다. 대괄호로 묶은 두 추출은 서로 기다리지 않고 함께 돌 수 있습니다. 보고서 태스크는 둘이 다 성공한 뒤에 돕니다.

flowchart TD
    A["extract_orders"] --> C["build_report"]
    B["extract_refunds"] --> C

Airflow 의 웹 화면도 DAG 를 이런 그림으로 보여 줍니다. 그림은 코드에서 나오므로 코드를 고치면 그림도 바뀝니다.

파이프라인이 코드라서 다른 코드처럼 저장소에 둡니다. 배포 전에는 리뷰를 거칩니다. 변경 이력이 남습니다. 모양이 비슷한 태스크 여러 개는 반복문으로 찍어 냅니다.

스케줄러와 워커

DAG 파일은 적어 두기만 한 것입니다. 이 파일을 읽고 때맞춰 실행하는 것은 Airflow 를 이루는 프로그램 몇 개입니다. 하나씩 보고 나서 한 그림으로 잇습니다.

스케줄러는 Airflow 의 중심입니다. DAG 파일을 되풀이해 읽어서 파이프라인 목록을 새로 고칩니다. 그리고 어느 태스크가 돌 때가 됐는지 따집니다. 정한 시각이 지났고 앞 태스크가 다 성공했으면 그 태스크를 실행하라고 넘깁니다.

스케줄러가 DAG 파일을 되풀이해 읽는다는 점에는 대가가 따릅니다. 파일 맨 위에 적은 코드는 읽을 때마다 돕니다. 거기에 데이터베이스 조회처럼 오래 걸리는 일을 넣으면 스케줄러 전체가 느려집니다. 그런 일은 태스크 안에 넣습니다.

메타데이터 데이터베이스는 Airflow 가 자기 상태를 적어 두는 관계형 데이터베이스입니다. 어느 DAG 가 언제 돌았고 태스크마다 어떻게 끝났는지가 여기에 쌓입니다. 스케줄러도 앞 태스크가 성공했는지를 여기서 읽고, 태스크를 넘겼다는 표시를 여기에 적습니다.

흔히 PostgreSQL 이나 MySQL 을 씁니다. 파이프라인이 옮기는 데이터가 담기는 곳은 아닙니다.

워커는 태스크를 받아 실행하는 프로세스입니다. 태스크의 코드를 돌리고 로그를 남깁니다. 끝나면 성공인지 실패인지를 메타데이터 데이터베이스에 적습니다.

실행기(executor)는 스케줄러가 넘긴 태스크를 어느 워커에서 돌릴지 정하는 부품입니다. 워커 프로세스를 어느 기계에 띄우느냐에 따라 실행기가 갈립니다.

실행기 가운데 하나는 작업 큐를 씁니다. 작업 큐는 할 일을 쌓아 두면 워커들이 하나씩 꺼내 가는 대기열입니다. Celery 는 파이썬에서 이런 작업 큐를 쓰게 해 주는 라이브러리입니다.

대표적인 실행기 셋을 견주면 아래와 같습니다.

실행기 태스크를 돌리는 곳
LocalExecutor 스케줄러와 같은 기계에 띄운 워커 프로세스
CeleryExecutor 여러 기계에 띄운 워커 프로세스가 작업 큐에서 나눠 가져감
KubernetesExecutor Kubernetes 에 태스크마다 컨테이너를 하나씩 띄움

LocalExecutor 는 기계 한 대로 끝나 가볍습니다. 나머지 둘은 태스크가 많아지면 기계를 늘려 받아 냅니다.

웹 화면은 이 데이터베이스를 읽어 보여 줍니다. DAG 마다 날짜별로 어느 태스크가 성공하고 실패했는지가 한 표에 뜹니다. 실패한 태스크를 골라 다시 돌리라고 누를 수도 있습니다.

지금까지 본 프로그램들을 이으면 아래와 같습니다. 스케줄러가 넘긴 태스크는 실행기를 거쳐 워커에서 돕니다. 그 결과는 모두 메타데이터 데이터베이스로 모입니다.

flowchart TD
    F["DAG 파일"] --> S["스케줄러"]
    S -->|상태를 읽고 쓴다| M[("메타데이터 데이터베이스")]
    S -->|때가 된 태스크| E["실행기"]
    E --> W["워커"]
    W -->|성공 · 실패 기록| M
    U["웹 화면"] -->|읽는다| M

태스크가 거치는 상태

태스크가 한 번 도는 것을 태스크 인스턴스(task instance)라고 부릅니다. 「build_report 태스크의 9월 24일 치 실행」이 태스크 인스턴스 하나입니다.

태스크 인스턴스는 몇 가지 상태를 거칩니다. 스케줄러가 돌 때가 됐다고 보면 예약됨(scheduled)이 됩니다. 실행기에 넘어가 워커를 기다리면 대기열(queued)입니다. 워커가 잡으면 실행 중(running)입니다. 끝나면 성공(success)이나 실패(failed)가 됩니다.

stateDiagram-v2
    state "예약됨" as scheduled
    state "대기열" as queued
    state "실행 중" as running
    state "재시도 대기" as retry
    state "성공" as success
    state "실패" as failed
    [*] --> scheduled
    scheduled --> queued
    queued --> running
    running --> success
    running --> retry: 실패 · 횟수 남음
    retry --> scheduled: 정한 시간 뒤
    running --> failed: 횟수를 다 씀
    success --> [*]
    failed --> [*]

재시도를 걸어 두면 실패가 곧바로 실패로 끝나지 않습니다. 재시도 대기(up_for_retry)로 갔다가 정한 시간이 지나면 다시 예약됩니다. 정한 횟수를 다 쓰고도 실패하면 그때 실패로 굳습니다.

앞 태스크가 실패하면 뒤 태스크는 돌지 않습니다. 이때 뒤 태스크는 상위 실패(upstream_failed) 상태가 됩니다. 덜 만들어진 데이터 위에서 뒤 단계가 도는 일이 이렇게 막힙니다.

데이터 품질 검사를 태스크로 넣는 방식이 이 규칙에 기댑니다. 적재 앞에 행 수나 빈 값을 세는 검사 태스크를 둡니다. 검사가 실패하는 날에는 뒤의 적재가 멈춥니다. 검사를 못 넘은 데이터가 보고서 표까지 가지 않습니다.

데이터 구간과 백필

Airflow 의 주기 실행에는 처음 보면 헷갈리는 규칙이 하나 있습니다. 이 절에서 실행은 DAG 전체가 한 번 도는 것을 가리킵니다. 실행 하나가 「언제 도느냐」와 「어느 날짜의 데이터를 맡느냐」가 다를 수 있습니다.

하루 한 번 도는 DAG 라면 실행 하나에 하루치 구간을 맡길 수 있습니다. 이 구간을 데이터 구간(data interval)이라고 부릅니다. 9월 24일 0시부터 25일 0시까지가 한 구간입니다.

구간을 맡은 실행은 구간이 끝난 뒤에 시작합니다. 24일 치 데이터는 24일이 다 지나야 모이기 때문입니다. 그래서 24일 치 실행은 25일 0시가 지나서 돕니다. 이 실행에는 구간의 시작인 24일이 이름으로 붙습니다. 이 날짜를 논리 날짜(logical date)라고 부릅니다.

맡는 구간 논리 날짜 도는 때
9월 23일 0시 ~ 24일 0시 9월 23일 9월 24일 0시 뒤
9월 24일 0시 ~ 25일 0시 9월 24일 9월 25일 0시 뒤

이 방식에서는 논리 날짜가 도는 날보다 하루 앞섭니다. 다른 방식은 구간을 두지 않고 도는 시각을 논리 날짜로 삼습니다. 이때는 25일 0시에 도는 실행의 논리 날짜가 25일입니다. 어느 쪽이 기본인지는 Airflow 버전마다 다르므로, 남의 DAG 를 읽을 때는 그 DAG 가 어느 방식으로 도는지부터 확인합니다.

어느 방식이든 태스크 코드는 오늘 날짜를 직접 읽지 않고 넘겨받은 날짜를 씁니다. 그래야 같은 실행을 언제 다시 돌려도 같은 날짜의 데이터를 처리합니다.

날짜를 넘겨받아 쓰면 지난 날짜를 다시 돌리기 쉬워집니다. 로직을 고친 뒤 지난 한 달 치 결과를 새 로직으로 다시 만드는 일이 백필입니다. Airflow 에 날짜 범위를 주면 그 안의 날짜마다 실행을 하나씩 만들어 돌립니다.

DAG 의 시작 날짜가 과거면 그 사이의 빈 구간을 어떻게 할지 정해야 합니다. 캐치업(catchup)을 켜 두면 스케줄러가 밀린 구간을 차례로 전부 돌립니다. 끄면 가장 최근 구간 하나만 돕니다. 앞 코드의 catchup=False 가 이 설정입니다.

같은 구간이 두 번 돌 수 있으니 태스크는 멱등하게 짭니다. 멱등하다는 것은 여러 번 돌려도 결과가 한 번 돌린 것과 같다는 뜻입니다. 적재 태스크가 그 날짜의 행을 먼저 지우고 다시 넣으면 됩니다. 재시도나 백필로 두 번 돌아도 행이 두 벌 쌓이지 않습니다.

태스크 사이에 오가는 값

Airflow 는 태스크의 순서를 챙길 뿐 데이터를 태스크 사이로 나르지 않습니다. 앞 예시의 orders.py 는 워커에서 돌며 주문을 꺼내고, 꺼낸 것을 파일 저장소에 내려놓고 끝나는 스크립트입니다. 다음 태스크는 거기서 읽습니다. 태스크끼리 데이터를 손에서 손으로 건네지 않습니다.

큰 계산은 아예 다른 시스템에 맡기기도 합니다. 대용량 데이터를 여러 기계로 나눠 계산하는 Apache Spark 나, 분석할 데이터를 한데 모아 두는 데이터 웨어하우스가 그런 시스템입니다. 이때 태스크는 계산을 시키는 명령만 보내고 끝나기를 기다립니다.

이렇게 다른 시스템에 일을 시키고 순서와 결과만 챙기는 일을 오케스트레이션이라고 부릅니다. Airflow 가 맡는 일이 바로 이것입니다.

작은 값은 태스크끼리 건넬 수 있습니다. 이 통로를 XCom(cross-communication, 태스크 사이 통신)이라고 부릅니다. 방금 만든 파일의 경로나 처리한 행 수 같은 짧은 값을 넘길 때 씁니다. 값이 메타데이터 데이터베이스에 저장되므로 큰 데이터를 싣는 통로가 아닙니다.

Airflow 가 맞지 않는 일

Airflow 는 주기를 두고 도는, 끝이 있는 일을 돌립니다. 한 번 돌면 끝나는 태스크를 날마다, 시간마다 다시 띄웁니다.

끝없이 흘러드는 데이터를 들어오는 즉시 처리하는 스트림 처리에는 맞지 않습니다. 그 일은 Apache Kafka 같은 메시지 플랫폼과 Apache Flink 같은 스트림 처리 엔진이 맡습니다. Airflow 는 쌓인 결과를 하루나 한 시간 단위로 모아 정리하는 일을 맡습니다.

태스크 하나를 띄우는 데에도 스케줄러, 실행기, 워커를 차례로 거칩니다. 이 과정이 태스크마다 되풀이됩니다. 그래서 몇 초 안에 답해야 하는 요청 처리에는 이 절차가 부담이 됩니다. 몇 초마다 도는 짧은 작업도 마찬가지입니다.

굴리는 비용도 있습니다. 스케줄러, 워커, 웹 화면, 메타데이터 데이터베이스를 띄우고 지켜봐야 합니다. 작업이 한두 개이고 서로 앞뒤가 없다면 cron 으로 충분합니다. 이 부담을 덜려고 클라우드 업체가 Airflow 를 대신 띄워 주는 관리형 서비스를 냅니다. Amazon MWAA 와 Google Cloud Composer 가 그런 서비스입니다.

관련 항목

Airflow 파이프라인을 이루는 구성 요소

DAG · 태스크 · 오퍼레이터 · 센서 · XCom · 데이터 구간 · 논리 날짜 · 캐치업

Airflow 를 굴리는 프로세스와 저장소

스케줄러 · 워커 · 메타데이터 · PostgreSQL · MySQL · Celery · Kubernetes · 작업 큐

Airflow 가 돌리는 데이터 작업

ETL · ELT · 배치 처리 · 백필 · 데이터 품질 · 파이프라인

Airflow 태스크가 일을 맡기는 처리 엔진

Apache Spark · dbt · 데이터 웨어하우스 · 데이터베이스

같은 일을 두고 겨루는 오케스트레이터

cron · Dagster · Prefect · Luigi · Argo Workflows · AWS Step Functions

Airflow 가 맡지 않는 실시간 처리

스트림 처리 · Apache Kafka · Apache Flink · 이벤트 스트림

Airflow 태스크가 지켜야 하는 성질

멱등성 · 재시도 · 원자성

Airflow 가 속하는 상위 분류

오케스트레이션 · 워크플로 · 워크플로우 오케스트레이션 · 데이터 엔지니어링

다른 이름: Airflow · 에어플로 · 에어플로우 · 아파치 에어플로