Local Executor
Local Executor
LocalExecutor가 스케줄러 노드에서 어떻게 프로세스를 생성해 태스크를 실행하는지 설명하는 문서예요. parallelism 파라미터의 의미와 기본값, Fork/Spawn 두 가지 생성 방식의 차이, 그리고 컨테이너 환경이나 다중 스케줄러 구성에서 주의할 점을 함께 살펴볼게요.
출처: 문서
본문
LocalExecutor는 스케줄러 노드에서 통제된 방식으로 프로세스를 생성해 태스크를 실행해요.
parallelism 파라미터는 노드를 과부하시키지 않도록 생성되는 프로세스 수를 제한해요. 이 파라미터는 반드시 0보다 커야 해요.
LocalExecutor는 start 시점에 self.parallelism 값과 같은 수의 프로세스를 생성해요. 이때 task_queue를 사용해 태스크 유입과 워커 간 작업 분배를 조정하며, 워커는 준비가 되는 대로 태스크를 가져가요. LocalExecutor의 수명 주기 동안 워커 프로세스는 태스크를 기다리며 실행되고 있다가, LocalExecutor가 셧다운 호출을 받으면 워커를 종료시키기 위한 poison token이 워커들로 보내져요.
워커 생성 동작은 multiprocessing 시작 방법에 따라 달라져요:
-
Fork mode (Linux 기본값): Copy-on-Write(COW)로 인한 메모리 급증을 막기 위해
parallelism까지의 워커를 한꺼번에 생성해요. 자세한 내용은 Discussion을 확인해 보세요. -
Spawn mode (macOS와 Windows 기본값): 많은 프로세스를 동시에 생성할 때의 오버헤드를 막기 위해 필요할 때마다 워커를 하나씩 생성해요.
참고 (Note)
parallelism파라미터는airflow.cfg의[core] parallelism옵션으로 설정할 수 있어요. 기본값은32예요.
경고 (Warning)
LocalExecutor 워커는 스케줄러의 하위 프로세스로 생성되기 때문에, 컨테이너 환경에서는 스케줄러 프로세스가 메모리를 과도하게 사용하는 것처럼 보일 수 있어요. 이는 OOM(Out of Memory)으로 인한 컨테이너 재시작을 유발할 수 있어요. 컨테이너의 리소스 한도에 맞춰
parallelism값을 조정하는 것을 고려해 보세요.
참고 (Note)
airflow.cfg의[core]섹션에서 여러 Scheduler를executor=LocalExecutor로 구성하면, 각 Scheduler가 각자 LocalExecutor를 실행해요. 즉 태스크가 스케줄러가 실행되는 머신들에 걸쳐 분산 방식으로 처리된다는 뜻이에요.한 가지 고려할 점이 있어요:
- Scheduler 재시작: Scheduler가 재시작되면, 다른 Scheduler들이 고아(orphaned)가 된 태스크를 인지하고 재시작하거나 실패 처리하는 데 시간이 걸릴 수 있어요.
참고 (Note)
이전 버전의 Airflow에는 LocalExecutor를 무제한 parallelism(
self.parallelism = 0)으로 구성하는 옵션이 있었어요. 이 옵션은 스케줄러 노드를 과부하시키지 않기 위해 Airflow 3.0.0에서 제거됐어요.