스케줄러
스케줄러 (Scheduler)
이 페이지는 Airflow Scheduler의 역할과 실행 방법을 설명해요. Scheduler는 모든 Task와 DAG를 모니터링하고 의존성이 완료된 Task 인스턴스를 트리거해요. airflow scheduler 명령으로 지속 서비스처럼 실행되며, 스케줄러를 여러 개 동시에 실행하는 HA(고가용성) 설정과 성능 미세 조정, 관련 설정 옵션을 다뤄요.
출처: 문서
본문
Airflow scheduler는 모든 Task와 DAG를 모니터링하고, 의존성이 완료되면 Task 인스턴스를 트리거해요. 백그라운드에서 scheduler는 지정된 Dag 디렉터리의 모든 DAG를 모니터링하고 동기화 상태를 유지하는 서브프로세스를 띄워요. 기본적으로 분당 한 번 scheduler는 Dag 파싱 결과를 수집하고 어떤 활성 Task를 트리거할 수 있는지 확인해요.
Airflow scheduler는 Airflow production 환경에서 지속 서비스(persistent service)로 실행되도록 설계됐어요. 시작하려면 airflow scheduler 명령만 실행하면 돼요. 이 명령은 airflow.cfg에 지정된 설정을 사용해요.
Scheduler는 실행할 준비가 된 Task를 실행하기 위해 설정된 Executor를 사용해요.
scheduler를 시작하려면 명령을 실행하기만 하면 돼요:
airflow scheduler
scheduler가 성공적으로 실행되면 DAG가 실행되기 시작해요.
Note
첫 번째 Dag Run은 DAG의 Task들의 최소
start_date를 기반으로 만들어져요. 이후 Dag Run은 DAG의 timetable에 따라 만들어져요.cron이나 timedelta 스케줄을 가진 DAG의 경우 scheduler는 그 기간이 끝날 때까지 Task를 트리거하지 않아요. 예를 들어
schedule이@daily로 설정된 작업은 하루가 끝난 후에 실행돼요. 이 기법은 해당 기간에 필요한 데이터가 DAG가 실행되기 전에 완전히 제공되도록 보장해요. UI에서는 Airflow가 Task를 하루 늦게 실행하는 것처럼 보여요.
Note
schedule이 하루인 DAG를 실행하면, 데이터 구간이2019-11-21에서 시작하는 run은2019-11-21T23:59이후에 트리거돼요.
다시 강조하면, scheduler는 start date보다 schedule 하나 뒤인, 구간의 END에서 작업을 실행해요.
Dag 스케줄링에 대한 자세한 내용은 Dag Runs를 참고해야 해요.
Note
Scheduler는 높은 처리량(high throughput)을 위해 설계됐어요. 이는 가능한 한 빨리 Task를 스케줄링하려는 정보에 입각한 설계 결정이에요. Scheduler는 pool에서 사용 가능한 빈 슬롯이 몇 개인지 확인하고 한 반복에서 그 수만큼의 Task 인스턴스를 최대로 스케줄링해요. 즉 Task 우선순위는 대기 중인 스케줄링된 Task가 큐 슬롯보다 많을 때만 발효돼요. 따라서 같은 배치를 공유하면 낮은 우선순위 Task가 높은 우선순위 Task보다 먼저 스케줄링되는 경우가 있을 수 있어요. 더 자세한 내용은 이 GitHub 토론을 참고해요.
스케줄러 여러 개 실행하기
Airflow는 성능과 복원력(resiliency)을 위해 스케줄러를 여러 개 동시에 실행하는 것을 지원해요.
개요
HA scheduler는 기존 메타데이터 데이터베이스를 활용하도록 설계됐어요. 이는 주로 운영 단순성을 위한 것이에요: 모든 컴포넌트가 이미 그 DB와 통신해야 하며, scheduler 사이의 직접 통신이나 합의 알고리즘(Raft, Paxos 등)이나 다른 합의 도구(Apache Zookeeper, Consul 등)를 사용하지 않음으로써 "운영 표면적"을 최소로 유지했어요.
스케줄러는 이제 직렬화된 Dag 표현을 사용해 스케줄링 결정을 내리며, 스케줄링 루프의 대략적인 개요는 다음과 같아요:
- 새 DagRun이 필요한 DAG가 있는지 확인하고 생성하기
- 스케줄 가능한 TaskInstance나 완료된 DagRun을 위해 DagRun 배치를 검사하기
- 스케줄 가능한 TaskInstance를 선택하고, Pool 제한과 다른 동시성 제한을 존중하면서 실행을 위해 큐에 넣기
다만 이는 데이터베이스에 몇 가지 요구사항을 둬요.
데이터베이스 요구사항
간단히 말하면 PostgreSQL 12+ 또는 MySQL 8.0+를 사용하는 사용자는 준비된 상태예요 — 원하는 만큼 스케줄러 사본을 실행할 수 있어요. 추가 설정이나 구성 옵션이 필요 없어요. 다른 데이터베이스를 사용한다면 계속 읽어보세요.
성능과 처리량을 유지하기 위해 스케줄링 루프의 한 부분이 메모리에서 많은 계산을 수행해요 (각 TaskInstance마다 DB로 왕복하면 너무 느리니까요). 따라서 한 번에 하나의 스케줄러만 이 임계 섹션(critical section)에 있어야 해요 — 그렇지 않으면 제한이 올바르게 존중되지 않을 거예요. 이를 위해 데이터베이스 행 레벨 잠금(SELECT ... FOR UPDATE 사용)을 사용해요.
이 임계 섹션은 TaskInstance가 scheduled 상태에서 executor로 큐에 들어가는 곳이에요. 그러면서 다양한 동시성·pool 제한이 존중되도록 보장해요. 임계 섹션은 Pool 테이블의 모든 행에 행 레벨 쓰기 잠금을 요청해 얻어요(대략 SELECT * FROM slot_pool FOR UPDATE NOWAIT와 같지만 정확한 쿼리는 약간 달라요).
다음 데이터베이스는 완전히 지원되며 "최적" 경험을 제공해요:
- PostgreSQL 12+
- MySQL 8.0+
Warning
MariaDB는 버전 10.6.0까지
SKIP LOCKED나NOWAITSQL 절을 구현하지 않았어요. 이 기능이 없으면 여러 스케줄러 실행은 지원되지 않으며 교착(deadlock) 에러가 보고됐어요. MariaDB 10.6.0 이후는 여러 스케줄러와 함께 올바르게 작동할 수 있지만, 테스트되진 않았어요.
Note
Microsoft SQL Server는 HA로 테스트되지 않았어요.
스케줄러 성능 미세 조정
스케줄러 성능에 영향을 주는 것들
Scheduler는 Task를 실행하기 위해 지속적으로 스케줄링할 책임이 있어요. scheduler를 미세 조정하려면 여러 요소를 고려해야 해요:
- 나의 배포 종류
- 사용 가능한 메모리가 얼마나 있는지
- 사용 가능한 CPU가 얼마나 있는지
- 사용 가능한 네트워킹 대역폭이 얼마나 있는지
- Dag 구조의 로직과 정의:
- DAG가 몇 개인지
- DAG가 얼마나 복잡한지 (Task·의존성이 몇 개인지)
- Scheduler 설정
- Scheduler가 몇 개인지
- Scheduler가 한 루프에서 처리하는 Task 인스턴스 수
- 루프마다 생성/스케줄링해야 할 새 Dag run 수
- Scheduler가 정리(cleanup)를 수행하고 고아(orphaned) Task를 확인/채택해야 하는 빈도
미세 조정을 수행하려면 Scheduler가 내부적으로 어떻게 작동하는지 이해하는 것이 좋아요. Airflow Summit 2021 강연 Deep Dive into the Airflow Scheduler talk를 보면 미세 조정을 수행하는 데 도움이 돼요.
Scheduler 미세 조정 접근법
Airflow는 성능을 미세 조정하기 위해 많은 "knob"을 제공하지만, 어떤 knob을 어떻게 돌려 최상의 효과를 얻을지는 사용자의 특정 배포, DAG 구조, 하드웨어 가용성, 기대치에 따라 달라지는 별개의 작업이에요. 배포를 관리하는 일의 일부는 무엇을 최적화할지 결정하는 것이에요.
Airflow는 결정할 유연성을 주지만, 어떤 성능 측면이 자신에게 가장 중요한지 파악하고 어떤 knob을 어느 방향으로 돌릴지 결정해야 해요.
일반적으로 미세 조정을 위한 접근 방식은 어떤 성능 개선·최적화와도 동일해야 해요 (특정 도구를 추천하진 않을게요 — 평소 시스템을 관찰·모니터링할 때 쓰는 도구를 쓰면 돼요):
- 평소 시스템을 모니터링할 때 쓰는 올바른 도구 세트로 시스템을 모니터링하는 것이 매우 중요해요. 이 문서는 사용할 수 있는 특정 메트릭·도구에 대해 자세히 다루지 않고, 어떤 종류의 자원을 모니터링해야 하는지만 설명해요. 올바른 데이터를 얻기 위해서는 모니터링에 대한 모범 사례를 따르면 돼요.
- 어떤 성능 측면이 자신에게 가장 중요한지 결정해요 (무엇을 개선하고 싶은지).
- 시스템을 관찰해 병목이 어디인지 확인해요: CPU, 메모리, I/O가 일반적인 제한 요소예요.
- 기대치와 관찰에 기반해 다음 개선이 무엇인지 결정하고, 다시 성능·병목 관찰로 돌아가요. 성능 개선은 반복적인 과정이에요.
Scheduler 성능을 제한할 수 있는 자원
주의해야 할 자원 사용 영역이 몇 가지 있어요:
- 성능을 높이고 더 많은 것을 병렬로 처리하려 할 때 데이터베이스 연결과 데이터베이스 사용이 문제가 될 수 있어요. Airflow는 "데이터베이스 연결을 많이 소비하는" 것으로 알려져 있어요 — DAG가 많을수록, 병렬로 많이 처리할수록 더 많은 데이터베이스 연결이 열려요. 이는 보통 MySQL에서는 문제가 되지 않는데, MySQL은 연결 처리가 스레드 기반이기 때문이에요. 하지만 Postgres에서는 문제가 될 수 있는데, 연결 처리가 프로세스 기반이기 때문이에요. 중간 규모의 Postgres 기반 Airflow 설치도 있다면 데이터베이스 프록시로 PGBouncer를 사용하는 것이 최선이라는 게 일반적인 합의예요. Apache Airflow용 Helm Chart는 PGBouncer를 기본적으로 지원해요.
- Airflow Scheduler는 여러 인스턴스와 함께 거의 선형적으로 확장되므로, Scheduler의 성능이 CPU 바운드라면 Scheduler를 더 추가할 수도 있어요.
- 메모리 사용량을 볼 때 관찰하는 메모리 종류에 주의하세요. 보통
total memory used보다는working memory를 봐야 해요 (이름은 배포에 따라 다를 수 있음).
Scheduler 성능을 개선하려면
자원 사용량을 알게 되면 고려할 수 있는 개선은 다음과 같아요:
- 자원 활용률을 개선하기. 시스템에 충분히 활용되지 않는 여유 용량이 있을 때(다시 말해 CPU, 메모리 I/O, 네트워킹이 주요 후보) — scheduler 수를 늘리거나 더 자주 작업하기 위해 간격을 줄이는 같은 조치가 더 높은 활용을 대가로 성능 개선을 가져올 수 있어요.
- 하드웨어 용량을 늘리기 (예: CPU가 제한으로 보일 때). 종종 scheduler 성능 문제는 단순히 시스템이 "충분히 유능"하지 않아서이며, 이것이 유일한 방법일 수 있어요. 예를 들어 머신의 모든 CPU를 사용하고 있다면 새 머신에 scheduler를 하나 더 추가하고 싶을 수 있어요 — 대부분의 경우 2번째·3번째 scheduler를 추가하면 스케줄링 용량이 선형적으로 늘어나요 (공유 데이터베이스 등이 병목이 아니라면).
- "scheduler 튜너블" 값으로 다양한 실험해 보기. 단순히 한 성능 측면을 다른 것과 교환해 더 나은 효과를 얻는 경우가 많아요. 보통 성능 튜닝은 서로 다른 측면을 균형 잡는 예술이에요.
Scheduler 설정 옵션
다음 설정을 사용해 Scheduler의 다양한 측면을 제어할 수 있어요. 그러나 [scheduler] 섹션에서 사용할 수 있는 다른 비-성능 관련 scheduler 설정 파라미터도 Configuration Reference에서 볼 수 있어요.
-
max_dagruns_to_create_per_loop
Dag run을 만들 때 각 scheduler가 잠그는 DAG 수를 변경해요. 이 값을 더 낮게 설정할 수 있는 한 가지 이유는 DAG가 엄청나게 크고(Dag당 10k+ Task 규모) 여러 scheduler를 실행 중일 때, 하나의 scheduler가 모든 작업을 하게 하고 싶지 않기 때문이에요.
-
max_dagruns_per_loop_to_schedule
Task를 스케줄링하고 큐에 넣을 때 scheduler가 검사(및 잠그)해야 할 DagRun의 수예요. 이 제한을 늘리면 작은 DAG에 대해 더 많은 처리량을 허용하지만, 큰 DAG(예: 500 Task 이상)의 처리량은 느려질 가능성이 있어요. 여러 scheduler를 사용할 때 이것을 너무 높게 설정하면 한 scheduler가 모든 Dag run을 차지해 다른 scheduler는 할 일이 없을 수 있어요.
-
Scheduler가 관련 쿼리에서
SELECT ... FOR UPDATE를 발행해야 하는지 여부. False로 설정하면 한 번에 하나 이상의 스케줄러를 실행하면 안 돼요. -
pool 사용 통계를 StatsD로 보내야 하는 빈도(초)(statsd_on이 활성화된 경우). 이를 계산하는 것은 비교적 비싼 쿼리이므로, StatsD 롤업 주기와 같은 주기로 설정해야 해요.
-
task instance(scheduled, queued, running, deferred) 통계를 StatsD로 보내야 하는 빈도(초)(statsd_on이 활성화된 경우). 이를 계산하는 것은 비교적 비싼 쿼리이므로, StatsD 롤업 주기와 같은 주기로 설정해야 해요.
-
Scheduler가 고아 Task나 죽은 SchedulerJobs를 확인해야 하는 빈도(초).
이 설정은 죽은 scheduler가 어떻게 감지되고 그것이 "감독"하던 Task를 다른 scheduler가 어떻게 가져가는지(pick up) 제어해요. Task는 계속 실행되므로, 잠시 감지하지 못해도 해는 없어요.
SchedulerJob이 "죽은" 것으로 감지되면(scheduler_health_check_threshold로 결정), 죽은 프로세스가 실행한 실행 중·대기 중 Task는 "채택(adopted)"되어 이 scheduler가 대신 모니터링해요.
-
max_tis_per_query 스케줄링 메인 루프의 쿼리 배치 크기예요.
core.parallelism보다 크면 안 돼요. 너무 높으면 쿼리 조건의 복잡성과/또는 과도한 잠금으로 SQL 쿼리 성능이 영향을 받을 수 있어요.또한 DB의 최대 허용 쿼리 길이에 도달할 수도 있어요.
core.parallelism의 값을 사용하려면 0으로 설정해요. -
scheduler_idle_sleep_time 루프에서 할 일이 없을 때 scheduler가 루프 사이에 얼마나 오래 sleep할지 제어해요. 즉 무언가를 스케줄링했다면 다음 루프 반복을 바로 시작해요. 이 파라미터는 이름이 좋지 않고(역사적 이유) 향후 현재 이름의 deprecation과 함께 개명될 거예요.