Перейти к содержанию
Learning Platform

Зачем вообще понимать планировщик

Большинство инженеров знают Airflow как «штуку, которая запускает DAG по расписанию». Этого достаточно ровно до того момента, пока всё работает. А потом задача висит в статусе scheduled десять минут вместо одной, или после падения сервера весь оркестратор встаёт колом, или вы запускаете два планировщика «для надёжности» и боитесь, что одна и та же задача выполнится дважды. Во всех этих ситуациях нужно понимать, что именно делает scheduler на каждой итерации своего цикла.

В этой статье мы разберём планировщик до уровня транзакций базы данных: что происходит в main loop, как DAG-файлы парсятся в отдельных процессах, что такое critical section и почему именно она делает безопасной работу нескольких активных планировщиков одновременно. Ориентируемся на Airflow 2.10/2.11 LTS, но архитектурные принципы остаются в силе и в 3.x.

Что планировщик делает на каждой итерации

Планировщик — это бесконечный цикл. Его сердце называется _run_scheduler_loop, и одна итерация этого цикла в общих чертах делает следующее:

  1. Обрабатывает результаты парсинга DAG-файлов (новые и изменённые DAG попадают в metadata database).
  2. Создаёт новые DagRun для тех DAG, у которых наступило время запуска по их timetable.
  3. Просматривает существующие DagRun в состоянии running и оценивает их task instances: какие зависимости выполнены, какие задачи стали runnable.
  4. Входит в critical section — внутри одной транзакции выбирает task instances, которые реально можно запустить прямо сейчас (с учётом пулов и лимитов параллелизма), и переводит их в состояние queued.
  5. Отдаёт поставленные в очередь задачи активному executor.
  6. Собирает «события» от executor (задача завершилась, упала, потерялась) и обновляет состояния в БД.

Ключевая идея: планировщик не выполняет код ваших задач. Он не запускает Python-функции внутри PythonOperator. Его работа — это менеджмент состояний в БД и принятие решения «вот эту задачу пора отправить исполнителю». Само исполнение происходит в worker-процессах executor, физически отдельных от планировщика.

Парсинг DAG: почему это отдельные процессы

Если бы планировщик парсил DAG-файлы прямо в своём главном цикле, любой тяжёлый или сломанный DAG-файл блокировал бы планирование всего кластера. Представьте DAG, который на верхнем уровне делает запрос в внешний API — каждый парсинг подвешивал бы весь Airflow.

Поэтому за парсинг отвечает отдельная подсистема — DagFileProcessor. Менеджер процессов (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-слота
runningworker / executorкод задачи реально исполняется
success / failedworker / executorзадача завершилась
up_for_retryпланировщикзадача упала, но остались попытки retry
up_for_reschedulesensor в режиме rescheduleсенсор освободил слот и ждёт следующей проверки
upstream_failed / skippedпланировщикзадача не будет запущена из-за состояния зависимостей

Граница ответственности проходит между queued и running. Всё, что слева — зона планировщика и базы данных. Переход в running и дальше — зона executor и worker.

Critical section: сердце планировщика

Теперь самое интересное. Когда планировщик собрал множество task instances в состоянии scheduled, он не может просто все их запустить. Нужно уважать ограничения:

  • parallelism — глобальный лимит одновременно работающих задач на весь кластер;
  • max_active_tasks_per_dag — лимит на конкретный DAG;
  • pool — задачи делят ограниченное число слотов.

Логика выбора живёт в методе _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;

Два хвостовых слова делают всю магию. FOR UPDATE ставит row-level lock на отобранные строки. А 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 и интерактивные песочницы. Большая часть материалов бесплатна и доступна без регистрации.

Ещё в направлении · Data Engineering

Все материалы направления →