Зачем вообще понимать планировщик
Большинство инженеров знают Airflow как «штуку, которая запускает DAG по расписанию». Этого достаточно ровно до того момента, пока всё работает. А потом задача висит в статусе scheduled десять минут вместо одной, или после падения сервера весь оркестратор встаёт колом, или вы запускаете два планировщика «для надёжности» и боитесь, что одна и та же задача выполнится дважды. Во всех этих ситуациях нужно понимать, что именно делает
В этой статье мы разберём планировщик до уровня транзакций базы данных: что происходит в main loop, как DAG-файлы парсятся в отдельных процессах, что такое critical section и почему именно она делает безопасной работу нескольких активных планировщиков одновременно. Ориентируемся на Airflow 2.10/2.11 LTS, но архитектурные принципы остаются в силе и в 3.x.
Что планировщик делает на каждой итерации
Планировщик — это бесконечный цикл. Его сердце называется _run_scheduler_loop, и одна итерация этого цикла в общих чертах делает следующее:
- Обрабатывает результаты парсинга DAG-файлов (новые и изменённые DAG попадают в
metadata databaseцентральную БД Airflow на PostgreSQL или MySQL, где хранятся DAG, DagRun, TaskInstance, Pool и состояние всего кластера ). - Создаёт новые
DagRunдля тех DAG, у которых наступило время запуска по их timetable. - Просматривает существующие
DagRunв состоянииrunningи оценивает их task instances: какие зависимости выполнены, какие задачи стали runnable. - Входит в critical section — внутри одной транзакции выбирает task instances, которые реально можно запустить прямо сейчас (с учётом пулов и лимитов параллелизма), и переводит их в состояние
queued. - Отдаёт поставленные в очередь задачи активному
executorкомпоненту, который физически запускает задачи: LocalExecutor, CeleryExecutor или KubernetesExecutor . - Собирает «события» от executor (задача завершилась, упала, потерялась) и обновляет состояния в БД.
Ключевая идея: планировщик не выполняет код ваших задач. Он не запускает Python-функции внутри PythonOperator. Его работа — это менеджмент состояний в БД и принятие решения «вот эту задачу пора отправить исполнителю». Само исполнение происходит в worker-процессах executor, физически отдельных от планировщика.
Парсинг DAG: почему это отдельные процессы
Если бы планировщик парсил DAG-файлы прямо в своём главном цикле, любой тяжёлый или сломанный DAG-файл блокировал бы планирование всего кластера. Представьте DAG, который на верхнем уровне делает запрос в внешний API — каждый парсинг подвешивал бы весь Airflow.
Поэтому за парсинг отвечает отдельная подсистема — DagFileProcessorManager) держит пул дочерних процессов и раздаёт им файлы из папки dags/ на парсинг. Каждый дочерний процесс импортирует Python-файл, находит объекты DAG, сериализует их структуру в JSON и пишет в таблицу serialized_dag.
Это даёт два важных свойства:
- Изоляция отказов. Сломанный DAG-файл роняет только свой дочерний процесс, а не планировщик целиком.
- Разделение чтения. Планировщик и веб-сервер читают DAG не из
.pyфайлов, а из сериализованного представления в БД. Поэтому изменение в коде DAG появляется в UI не мгновенно, а после следующего цикла парсинга — это регулируется параметромmin_file_process_interval.
Цикл планирования и цикл парсинга работают независимо. Частота главного цикла регулируется параметром scheduler_idle_sleep_time (раньше назывался processor_poll_interval): когда делать нечего, планировщик «спит» долю секунды, чтобы не жечь CPU вхолостую. Под нагрузкой же он крутится практически без пауз.
DAG, TaskInstance и их состояния
Чтобы говорить о critical section предметно, нужно держать в голове модель состояний. Вот ключевые состояния TaskInstance, через которые задача проходит под управлением планировщика:
| Состояние | Кто выставляет | Что означает |
|---|---|---|
none / null | планировщик | task instance создан, но зависимости ещё не оценены |
scheduled | планировщик | зависимости выполнены, задача готова к запуску, ждёт ресурсов |
queued | планировщик (в critical section) | задача отдана executor, ждёт свободного worker-слота |
running | worker / executor | код задачи реально исполняется |
success / failed | worker / executor | задача завершилась |
up_for_retry | планировщик | задача упала, но остались попытки retry |
up_for_reschedule | sensor в режиме reschedule | сенсор освободил слот и ждёт следующей проверки |
upstream_failed / skipped | планировщик | задача не будет запущена из-за состояния зависимостей |
Граница ответственности проходит между queued и running. Всё, что слева — зона планировщика и базы данных. Переход в running и дальше — зона executor и worker.
Critical section: сердце планировщика
Теперь самое интересное. Когда планировщик собрал множество task instances в состоянии scheduled, он не может просто все их запустить. Нужно уважать ограничения:
parallelism— глобальный лимит одновременно работающих задач на весь кластер;max_active_tasks_per_dag— лимит на конкретный DAG;poolименованный набор слотов в Airflow, ограничивающий число одновременных задач, конкурирующих за общий ресурс, например за подключения к внешней БД — задачи делят ограниченное число слотов.
Логика выбора живёт в методе _critical_section_enqueue_task_instances, и называется она critical section не случайно. Это секция, которую в любой момент времени должен исполнять только один планировщик в отношении конкретного набора строк. Вся она обёрнута в одну транзакцию БД.
Внутри транзакции планировщик делает запрос примерно такого вида (упрощённо):
SELECT *
FROM task_instance
WHERE state = 'scheduled'
ORDER BY priority_weight DESC, execution_date ASC
LIMIT :max_tis_per_query
FOR UPDATE SKIP LOCKED;
Два хвостовых слова делают всю магию. SKIP LOCKED означает «не жди освобождения заблокированных строк, просто пропусти их и возьми следующие свободные».
Получив набор кандидатов, планировщик в памяти прогоняет их через лимиты: проверяет, есть ли свободные слоты в пуле, не превышен ли parallelism, не упёрся ли DAG в max_active_tasks_per_dag. Те задачи, что проходят все проверки, он переводит в queued и фиксирует транзакцию (COMMIT). На этом строковые блокировки снимаются. Задачи, которые не прошли по лимитам, остаются в scheduled до следующей итерации.
from datetime import datetime
from airflow import DAG
from airflow.operators.python import PythonOperator
with DAG(
dag_id="critical_section_demo",
schedule="@hourly",
start_date=datetime(2026, 1, 1),
catchup=False,
max_active_tasks=4, # потолок параллелизма для этого DAG
) as dag:
def pull_from_api(**context):
# эта функция исполняется на worker, НЕ в планировщике
...
extract = PythonOperator(
task_id="extract",
python_callable=pull_from_api,
pool="external_api", # задача делит слоты пула external_api
priority_weight=10, # выше приоритет в ORDER BY critical section
)
transform = PythonOperator(task_id="transform", python_callable=lambda: None)
extract >> transform
В этом примере priority_weight напрямую влияет на ORDER BY в SQL-запросе critical section: при нехватке слотов планировщик предпочтёт задачи с большим весом. А pool="external_api" гарантирует, что число одновременных обращений к внешнему API не превысит размер пула — даже если десятки DAG захотят дёрнуть его одновременно.
High Availability: почему два планировщика не запустят задачу дважды
До Airflow 2.0 планировщик был единой точкой отказа: его можно было запускать только в одном экземпляре. Падает планировщик — встаёт весь оркестратор. В Airflow 2.0 появилась поддержка нескольких active-active планировщиков, и сделана она элегантно: никакого внешнего сервиса для выбора лидера, никакого ZooKeeper, никакого Raft. Координация полностью построена на той самой critical section.
Логика такая. Запустите хоть три планировщика на трёх машинах — они все подключены к одной metadata database. Каждый на своей итерации заходит в critical section и выполняет тот самый SELECT ... FOR UPDATE SKIP LOCKED. И вот что происходит:
- Планировщик A берёт блокировку на строки задач
[1, 2, 3]. - Планировщик B в тот же момент выполняет тот же запрос. Строки
[1, 2, 3]для него заблокированы — благодаряSKIP LOCKEDон их пропускает и берёт[4, 5, 6]. - Оба планировщика работают с непересекающимися наборами задач, ставят их в
queuedи фиксируют транзакции.
Никакая задача не может быть отобрана двумя планировщиками одновременно, потому что row-level lock физически не позволит второй транзакции увидеть уже захваченную строку как свободную. SKIP LOCKED превращает потенциальную блокировку-ожидание в безопасное разделение работы: планировщики не конкурируют за одни и те же строки, а естественным образом «расходятся» по разным.
Именно поэтому в проде Airflow HA выглядит так буднично: вы просто запускаете несколько процессов airflow scheduler, и они балансируют нагрузку и страхуют друг друга без какой-либо явной координации. Цена этого решения — требование к СУБД: SELECT ... FOR UPDATE SKIP LOCKED должен поддерживаться. PostgreSQL поддерживает его с версии 9.5, MySQL — с 8.0. Поэтому официально для HA нужен один из них; старые версии MySQL и SQLite не подходят.
Несколько практических следствий HA-режима:
- Без single point of failure. Падение одного планировщика означает лишь временное снижение пропускной способности, а не остановку кластера.
- Горизонтальное масштабирование планирования. При тысячах DAG несколько планировщиков делят между собой работу по оценке task instances.
- Нагрузка на БД растёт. Каждый планировщик постоянно держит транзакции и берёт блокировки, поэтому metadata database становится критическим компонентом — её latency напрямую определяет потолок планирования.
Executor handoff и почему важна latency планировщика
Когда задача стала queued, планировщик передаёт её executor. Дальше путь зависит от типа исполнителя:
- LocalExecutor порождает дочерний процесс прямо на хосте планировщика.
- CeleryExecutor кладёт сообщение в брокер (Redis или RabbitMQ), откуда его забирает свободный Celery-worker.
- KubernetesExecutor создаёт отдельный pod на каждую задачу.
Worker, получив задачу, переводит её в running и запускает airflow tasks run, который и исполняет код оператора. Здесь важна одна деталь: между queued и running есть «слепая зона». Если worker по какой-то причине не подхватил задачу (упал, не хватило ресурсов), задача может зависнуть в queued. Для этого у планировщика есть таймаут task_adoption_timeout и механизм, который «усыновляет» или сбрасывает зависшие задачи обратно в scheduled.
Теперь о латентности. Между «задача стала runnable» и «задача реально пошла» проходит как минимум одна итерация главного цикла планировщика. Если цикл занимает, скажем, 5 секунд (тяжёлый парсинг, медленная БД, тысячи task instances для оценки), то каждая задача в среднем ждёт лишних несколько секунд. Для DAG с сотнями коротких задач это превращается в заметный оверхед: суммарное время выполнения DAG растёт не из-за самих задач, а из-за времени ожидания планировщика между ними.
Поэтому в production важно:
- держать metadata database быстрой (это узкое место critical section);
- не делать тяжёлых операций на верхнем уровне DAG-файлов (это замедляет парсинг);
- настраивать
max_tis_per_queryи число планировщиков под свою нагрузку; - мониторить метрику
scheduler.scheduler_loop_duration— она показывает, сколько реально длится одна итерация.
Соберём картину целиком
Планировщик Airflow — это не «крон с UI», а распределённый менеджер состояний поверх реляционной БД. Его main loop парсит DAG в изолированных процессах, создаёт DagRun по timetable, оценивает зависимости task instances и в рамках одной транзакции — critical section — выбирает, что запустить, уважая пулы и лимиты параллелизма. Та же самая транзакция с SELECT ... FOR UPDATE SKIP LOCKED бесплатно решает задачу High Availability: несколько активных планировщиков безопасно делят работу, и ни одна задача не уходит на исполнение дважды. А скорость одной итерации этого цикла напрямую определяет, как быстро ваши задачи переходят от scheduled к running.
Если хочется не просто прочитать про это, а разобрать main loop по исходникам, потрогать critical section в песочнице и собрать настоящий HA-сетап — на нашей платформе есть бесплатный курс по Apache Airflow 2.x с модулем «Scheduler internals», где всё это разбирается с интерактивными диаграммами и до уровня транзакций. Курс полностью открыт, без оплаты и регистрации.
Что дальше
Scheduler — лишь одна из тем в инженерии данных. Если вы строите карьеру в этом направлении, посмотрите весь каталог направления Data Engineering: оркестрация, потоковая обработка, хранилища, форматы данных — всё с уклоном в internals и интерактивные песочницы. Большая часть материалов бесплатна и доступна без регистрации.