
앞선 글에서는 Airflow에서 Workflow를 DAG로 정의하는 방법을 살펴봤다.
이번 글은 정의된 DAG가 실제로 어떻게 실행되는지를 Scheduler와 Executor 중심으로 정리한다.
DAG는 Workflow의 정의다.
정의된 DAG가 실제로 실행되면 DAG Run이 생성되고, Scheduler가 실행 가능한 Task를 판단한다. 이후 Executor를 통해 Task가 실행된다.
DAG 정의
↓
DAG Run 생성
↓
Scheduler
↓
실행 가능한 Task 판단
↓
Executor
↓
Task 실행
이 과정에서 Scheduler와 Executor의 역할을 구분해보면 다음과 같다.
Scheduler는 DAG와 Task의 상태를 계속 확인하면서 실행 가능한 Task Instance를 찾아 실행을 스케줄링하는 역할을 한다.
다음과 같은 Workflow가 있다고 가정한다.
extract
↓
transform
↓
load
extract가 완료되어 transform을 실행할 수 있는 상태가 되면 Scheduler가 해당 상태를 확인하고 실행 대상으로 판단한다.
즉, Scheduler는 단순히 정해진 시간에 Workflow를 시작하는 역할만 하는 것이 아니다.
현재 어떤 Task를 실행할 수 있는지 판단하고, 실행할 Task를 Executor에 전달하는 것 까지를 수행한다.
Scheduler가 실행할 Task를 판단했다면, 실제 Task 실행은 Executor가 담당한다.
Scheduler
│
│ 실행할 Task 결정
↓
Executor
│
│ 실행 방식에 따라 전달
↓
실행 환경
두 구성 요소의 역할을 정리하면 다음과 같다.
| 구성 요소 | 역할 |
|---|---|
| Scheduler | 실행 가능한 Task를 판단하고 스케줄링 |
| Executor | Task를 실행할 방식 결정 |
| Worker / Pod 등 | 실제 Task가 실행되는 환경 |
Airflow 3에서는 Executor가 별도의 독립 서버나 프로세스로 동작하는 구조가 아니다.
Executor는 Scheduler 내부에서 Task 실행을 처리하는 메커니즘으로 동작한다.
따라서 다음과 같이 이해하는 편이 정확하다.
Scheduler
├── Workflow Scheduling
└── Executor
↓
실행 환경
↓
Task 실행
Executor는 Task Instance를 실제로 실행하기 위한 실행 메커니즘이다.
Airflow는 Executor에 따라 Task를 실행하는 방식과 환경이 달라진다.
대표적으로 다음 세 가지 Executor를 볼 수 있다.
LocalExecutor는 Scheduler가 실행되는 환경에서 Task를 로컬로 실행한다.
Scheduler
│
└── LocalExecutor
↓
Local Task 실행
별도의 원격 Worker가 필요하지 않기 때문에 비교적 단순한 환경 & 분산 실행이 필요없는 경우에 적합하다.
CeleryExecutor는 여러 Worker에 Task를 분산해 실행한다.
Scheduler
│
CeleryExecutor
│
┌─────────┼───────────┐
↓ ↓ ↓
Worker 1 Worker 2 Worker 3
│ │ │
Task Task Task
여러 Worker가 Task를 나눠 실행하므로 다수의 Task를 분산 처리 (동시처리) 할 수 있다.
KubernetesExecutor는 Task 실행 시점마다 Kubernetes Pod를 생성해 그 안에서 Task를 실행한다.
Scheduler
│
KubernetesExecutor
│
┌──────────┼──────────┐
↓ ↓ ↓
Pod A Pod B Pod C
Task A Task B Task C
Task마다 독립된 Pod를 사용하기 때문에, Task별로 실행 환경과 리소스를 완전히 분리할 수 있다.
서로 다른 Task가 서로 다른 의존성이나 리소스 사양을 필요로 하는 경우 특히 유용하다.
Airflow
→ Workflow의 순서와 실행 상태를 관리함
Kubernetes
→ Container / Pod의 실행 환경을 제공함
: Task를 어떤 실행 환경에서 처리할 것인지 결정
Airflow
│
Scheduler
│
Executor
│
┌───────────┼─────────────┐
↓ ↓ ↓
Local Worker Pod
LocalExecutor Celery Kubernetes
LocalExecutor
→ 로컬에서 실행CeleryExecutor
→ 여러 Worker에 분산하여 실행KubernetesExecutor
→ Task마다 독립된 Kubernetes Pod에서 실행
앞에서 살펴본 내용을 extract → transform → load Workflow에 적용해보자.
Scheduler
↓
extract 실행 가능 여부 확인
↓
extract 실행 스케줄링
↓
extract 실행
↓
transform 실행 가능 여부 확인
↓
transform 실행 스케줄링
↓
transform 실행
↓
load 실행 가능 여부 확인
↓
load 실행 스케줄링
↓
load 실행
여기서 중요한 점은 Scheduler가 실행 순서를 새로 만드는 것은 아니라는 것이다.
extract → transform → load라는 관계는 이미 DAG에 정의되어 있다.
Scheduler는 DAG에 정의된 관계와 각 Task의 상태를 확인하면서 현재 실행할 수 있는 Task를 스케줄링한다.
그 이후 Executor가 해당 Task를 설정된 방식에 따라 실행한다.
| 개념 | 핵심 역할 |
|---|---|
| DAG | Workflow 정의 |
| DAG Run | Workflow의 실제 실행 단위 |
| Scheduler | 실행할 Task 판단 및 스케줄링 |
| Dependency | Task 실행 순서 정의 |
| Executor | Task 실행 방식 결정 |
| Worker / Pod | 실제 Task 실행 환경 |
특히 이번 글에서 기억할 부분은 두 가지다.
- Scheduler는 실행할 Task를 판단한다.
- Executor는 그 Task를 어떤 방식으로 실행할지 결정한다.
Executor에 따라 실제 실행 환경도 달라진다.
| 구성 요소 | 핵심 질문 |
|---|---|
| Scheduler | 지금 어떤 Task를 실행할 수 있는가? |
| Executor | 그 Task를 어떤 방식으로 실행할 것인가? |
| 실행 환경 | 실제로 어디에서 실행되는가? |
이번 글에서는 DAG로 정의한 Workflow가 실제로 실행되는 과정에서 Scheduler와 Executor가 각각 어떤 역할을 맡는지 정리했다.
Airflow의 Workflow 실행 구조
DAG
↓
DAG Run
↓
Scheduler
↓
Executor
↓
실행 환경
↓
Task 실행
Executor에 따른 실행 환경
Scheduler
│
└── Executor
│
├── LocalExecutor
│ └── Local
│
├── CeleryExecutor
│ └── Worker
│
└── KubernetesExecutor
└── Pod
다음 글에서는 실행 방식에서 한 단계 더 나아가,
실제 Workflow를 어떻게 설계할 것인지 다룰 예정이다.