Проблема з cron
До того, як інструменти оркестрації стали стандартом, типовий 'конвеєр' даних складався з папки з завданнями cron, кожне з яких запускалося в певний час і сподівалося, що попереднє завдання вже завершилось. Це працювало доти поки що не стало не так: завдання, яке виконувалось на дванадцять хвилин пізніше через перевантаження складу, безшумно пошкоджувало вихідні дані всього, що планувалось після нього, і ніхто не помічав цього, поки аналітик не запитає, чому цифри дашборду минулого дня виглядають неправильно. Cron кодує лише час, ніколи не визначає залежності, тому він структурно не може виражати 'завдання B має виконуватися після успішного виконання завдання A, а не просто з певним часовим зсувом'.
Основний крок Airflow полягав у заміні тригерів на основі часу діаграмою залежностей. Конвеєр виражається як направлений циклічний граф (DAG): набір завдань (вузлів), з’єднаних ребрами, які означають 'це завдання має завершитися перед початком іншого'. Оскільки граф не є циклічним, немає можливості випадково створити нескінченний цикл залежностей — Airflow перевіряє це під час парсингу та відмовляється розкладати DAG, який не є правильним DAG. Робоча роль планувальника — обхід графа: проходиться по DAG, знаходить завдання, для яких усі попередні залежності успішно виконані, і передає їх робочому процесу.
Дати виконання та ідентичні запуски
Один із більш тонких, але важливих дизайнерських виборів Apache Airflow полягає в тому, що кожен запуск DAG позначається логічною датою виконання, яка представляє період даних, який він призначений обробляти. Ця дата не збігається з фактичним часом виконання запуску, який був здійснений. DAG, запланований для запуску щодня о 2:00 ранку для ‘2026-07-24’, обробляє дані 24-го числа, навіть якщо його було запущено пізно або перепущено три дні після фактичного терміну – логічна дата, а не годинник, визначає, з якими даними працює задача. Це роз’єднання робить можливим
заповнення
>– якщо ви змінюєте трансформацію і потрібно повторно обробити останні 90 днів, ви просто запускаєте 90 DAG-ів, по одному на історичну логічну дату, і кожен з них запитує розділ за свою дату замість ‘сьогодні’.
Для цього необхідно, щоб завдання були ідентичними – повторне виконання однієї й тієї ж логічної дати повинно давати однаковий результат, а не подвоєння рядків. Airflow не забезпечує ідентичність для вас; це контракт, який повинен дотримуватися автор конвеєра, зазвичай шляхом того, щоб кожне завдання перезаписувало конкретний розділ (`INSERT OVERWRITE PARTITION` або еквівалентно), а не бездумно додавало дані. Команди, які ігнорують цю дисципліну, дізнаються про вартість лише тоді, коли тимчасова помилка змушує задачу повторно виконати її та їх таблиця доходів непомітно подвоюється.
Оператори, датчики та абстракція завдань
Вузли DAG будуються з операторів, які є параметризованими одиницями роботи — BashOperator виконує команду shell, PythonOperator виконує callable на Python, SnowflakeOperator виконує SQL-запит проти сховища, і так далі. Це свідомо загальний підхід: Airflow не піклується про те, що робить завдання, лише коли воно має виконуватися відносно сусідніх завдань та яке його політика повторних спроб. Кожне завдання незалежно може вказувати кількість повторів, затримку між ними (зазвичай з експоненційним зменшенням), тайм-аут і хуки сповіщення, тому непередбачуваний виклик API, який зазнає невдачі лише один раз на двадцять спроб, може автоматично перероблятися без зриву всієї конвеєра або пробудження когось.
Датчики – це особливий тип операторів, які не виконують роботу, а чекають на певну умову — FileSensor опрацьовує файл до появи в хмарному сховищі, ExternalTaskSensor чекає завершення завдання в іншій DAG, а датчик розділу чекає, поки дані за певний день з’являться в таблиці джерела. Датчики дозволяють виражати міжсистемні залежності ('не починайте цю трансформацію, поки робота з надсиланням даних трьома командами далі не завершиться сьогодні') без жорсткого кодування фіксованого часу очікування, що є тією ж категорією помилки, яку заохочував cron у перше місце. Сучасні версії Airflow також підтримують 'відкладаючі' датчики, які звільняють свій робочий слот під час очікування, тому флот конвеєрів, що чекають на повільні upstream-системи, не позбавляє планувальника ресурсів.
Цикл планувальника та де саме виконуються завдання
Під капотом, планувальник Airflow безперервно аналізує ваші файли визначення DAG (Python-скрипти, які створюють об'єкти графа), оцінює, які інстанції завдань мають право на виконання на основі їхніх залежностей та розкладу DAG, і додає відповідні завдання до черги. Віддільний виконавчий механізм (executor) вирішує, як саме виконуються ці завдання: найпростіший SequentialExecutor виконує одне завдання за раз і підходить лише для локального тестування; CeleryExecutor та KubernetesExecutor відповідно розподіляють завдання між пулом робітників або створюють новий pod на кожне завдання, що використовується в реальних розгортаннях. Цей розділ — планувальник вирішує, які завдання потрібно виконати, а виконавець визначає, де саме їх виконують — відображає ту ж саму стратегію «розділяй та володарюй», що й у поділі ResourceManager та NodeManager в YARN.
Візуально DAG, відтворений в інтерфейсі Airflow UI, виглядає як виробничий цех: коробки для кожного завдання, кольоровізовані за станом (чекають, виконуються, успішно, невдало, помилка попереднього завдання), з'єднані стрілками, що показують потік залежностей, а також Gantt-відображення, яке точно відображає час виконання кожного завдання відносно сусідніх. Це візуалізація не випадкова — це основний інструмент налагодження, оскільки коли щось ламається о 3 годині ночі, перше питання завжди таке: «Який вузол у графі став червоним і чи поширилася його помилка на нижче розташовані завдання?» DAG-відображення відповідає на це питання одразу.
Часті запитання
Що таке 'DAG' у Airflow?
Це абревіатура Directed Acyclic Graph – набір завдань, з’єднаних напрямними залежностями, без циклів, що гарантує завжди можливий порядок їх виконання.
Чому ідемпотентність така важлива для завдань Airflow?
Оскільки Airflow повторно запускає невдалі завдання та підтримує перезапуск історичних дат для бекфілів. Завдання, яке не є ідемпотентним – наприклад, завдання, що додає дані замість перезапису розділу – може безшумно генерувати дублікати або неправильні дані при будь-якому повторному запуску чи перезапуску.
Яка різниця між планувальником та виконавцем?
Планувальник аналізує DAGs і вирішує, які інстанції завдань готові до виконання на основі залежностей та розкладу; виконавець – це окремий компонент, який визначає, як і де фактично виконуються готові завдання, чи локально, через пул Celery worker або як Kubernetes pods.
Як бекфіл відрізняється від звичайного запланованого запуску?
Бекфіл запускає DAG runs для діапазону логічних дат виконання минулого періоду замість очікування, поки розклад їх досягне, дозволяючи перепрацювати історичні дані з оновленою конвеєризацією без очікування реальних днів, поки планувальник «дожене» запізнення.
Спробуйте наживо
Усе, що вище, працює прямо у вашому браузері — відкрийте the simulation і змінюйте параметри під час роботи. Нічого не встановлюється, нічого не завантажується на сервер, уся модель живе в одній вкладці.
▶ Відкрити симуляцію the simulation