The Mental ModelМентальная модель
▾Airflow is an orchestrator, not a compute engine. It schedules, sequences, retries, and tracks state of tasks. It does not process big data itself — it triggers Spark/dbt/SQL/containers that do.
Airflow — это оркестратор, а не движок вычислений. Он планирует, упорядочивает, перезапускает и отслеживает состояние задач. Он не обрабатывает большие данные сам — он запускает Spark/dbt/SQL/контейнеры, которые это делают.
Say this out loud in the call: "Airflow gives me at-least-once execution; correct outputs come from making tasks idempotent." That single sentence signals seniority.
Произнеси это вслух на созвоне: «Airflow даёт мне at-least-once execution; корректные результаты получаются за счёт того, что задачи идемпотентны». Одно это предложение — сигнал сеньорности.
End-to-End Run LifecycleЖизненный цикл выполнения
▾@daily run for June 13 actually fire?" →
at the END of the interval (00:00 June 14). logical_date = start of the interval, lagging wall-clock by one interval. Get this right or lose credibility.
@daily за 13 июня?» →
в КОНЦЕ интервала (00:00 14 июня). logical_date = начало интервала, отстаёт от реального времени на один интервал. Если ошибёшься — потеряешь доверие.
Core Objects & VocabularyБазовые объекты и терминология
▾| Term | What it is |
|---|---|
| DAG | Directed Acyclic Graph: the pipeline definition. Acyclic ⇒ valid topological order, no infinite loops. |
| DagRun | One execution of a DAG for a specific data interval / logical date. |
| Operator | Template/class for a unit of work (PythonOperator, BashOperator, SQL ops, KubernetesPodOperator). |
| Task | An operator instantiated in a DAG (a node in the graph). |
| TaskInstance | One run of one task for one DagRun — the thing that actually has state. |
| Sensor | Special operator that waits for a condition (file, partition, external task). |
| XCom | Cross-communication: small key/value passed between tasks via metadata DB. |
| Connection / Variable | Credentials+endpoints / arbitrary global config k-v. |
| Pool | Caps concurrent task slots for a scarce resource (e.g. a fragile DB). |
| Triggerer | Async process running deferrable operators so waiting tasks don't hold worker slots. |
| Термин | Что это |
|---|---|
| DAG | Directed Acyclic Graph: определение пайплайна. Acyclic ⇒ валидный топологический порядок, нет бесконечных циклов. |
| DagRun | Одно выполнение DAG для конкретного data interval / logical date. |
| Operator | Шаблон/класс для единицы работы (PythonOperator, BashOperator, SQL ops, KubernetesPodOperator). |
| Task | Оператор, проинстанцированный в DAG (узел в графе). |
| TaskInstance | Один запуск одной задачи для одного DagRun — вот он-то и имеет состояние. |
| Sensor | Специальный оператор, который ждёт условия (файл, партиция, внешняя задача). |
| XCom | Cross-communication: маленький key/value, передаваемый между задачами через metadata DB. |
| Connection / Variable | Credentials+endpoints / произвольный глобальный конфиг k-v. |
| Pool | Ограничивает одновременные слоты задач для дефицитного ресурса (напр. хрупкой БД). |
| Triggerer | Асинхронный процесс для deferrable операторов, чтобы ожидающие задачи не занимали воркер-слоты. |
Analogy: Operator : Task : TaskInstance ≈ class : object : invocation-with-state.
Аналогия: Operator : Task : TaskInstance ≈ класс : объект : вызов-с-состоянием.
# TaskFlow API (modern, preferred over raw operators) from airflow.decorators import dag, task @dag(schedule="@daily", start_date=datetime(2026,1,1), catchup=False) def pipeline(): @task def extract(ds=None): # ds = logical date string return f"s3://raw/{ds}/file.parquet" # small pointer → XCom OK @task def load(path): ... load(extract()) # dependency inferred from data flow pipeline()
# TaskFlow API (современный, предпочтительнее сырых операторов) from airflow.decorators import dag, task @dag(schedule="@daily", start_date=datetime(2026,1,1), catchup=False) def pipeline(): @task def extract(ds=None): # ds = строка логической даты return f"s3://raw/{ds}/file.parquet" # маленький указатель → XCom ОК @task def load(path): ... load(extract()) # зависимость выводится из потока данных pipeline()
Scheduling, Dependencies & Trigger RulesРасписание, зависимости и правила триггеров
▾Scheduling knobs key
Настройки расписания ключевое
start_date— first interval considered. Never dynamic (datetime.now()) → undefined first interval.schedule— cron / preset (@daily) / timedelta /None/ dataset list.catchup—Truebackfills every interval from start_date on unpause (can spawn 100s of runs). Default tocatchup=False+ controlled backfills.max_active_runs/max_active_tasks/parallelism— concurrency blast-radius control.
start_date— первый рассматриваемый интервал. Никогда динамический (datetime.now()) → неопределённый первый интервал.schedule— cron / пресет (@daily) / timedelta /None/ список dataset.catchup—Trueзаполняет каждый интервал от start_date при снятии с паузы (может породить сотни запусков). По умолчанию ставьcatchup=False+ контролируемые backfills.max_active_runs/max_active_tasks/parallelism— контроль радиуса поражения параллелизма.
Dependencies & Trigger Rules
Зависимости и правила триггеров
extract >> transform >> load # or chain(), set_downstream()# или chain(), set_downstream()
| Trigger rule | Fires when… | Use for |
|---|---|---|
all_success (default) | all upstream succeeded | normal flow |
all_done | all upstream finished (any state) | cleanup / notify regardless |
none_failed_min_one_success | none failed, ≥1 succeeded | join after branching |
one_failed / one_success | ≥1 in that state | fast-fail / fan-in alerts |
| Правило триггера | Срабатывает когда… | Для чего |
|---|---|---|
all_success (по умолч.) | все upstream успешны | обычный поток |
all_done | все upstream завершены (любое состояние) | cleanup / уведомление в любом случае |
none_failed_min_one_success | ни один не упал, ≥1 успешен | join после ветвления |
one_failed / one_success | ≥1 в этом состоянии | быстрый fail / fan-in алерты |
trigger_rule="all_done".
"Conditional path?" → @task.branch + none_failed_min_one_success on the join.
trigger_rule="all_done".
«Условный путь?» → @task.branch + none_failed_min_one_success на join.
Idempotency, Backfills & RetriesИдемпотентность, backfills и ретраи most-askedчастый вопрос
▾Idempotency = same logical date, run N times → same end state core
Идемпотентность = та же logical date, запущено N раз → то же конечное состояние суть
Required because retries, manual reruns, and backfills will re-execute tasks.
Требуется, потому что ретраи, ручные перезапуски и backfills неизбежно переисполнят задачи.
- Partition-scoped writes — write only the partition for
{{ ds }}, never "everything since yesterday." - Overwrite-by-partition / MERGE / upsert — a rerun replaces, never appends. (dbt
--full-refreshon a partition;INSERT OVERWRITE; IcebergreplaceWhere/MERGE.) - Deterministic inputs — parameterize off
logical_date, neverdatetime.now()/CURRENT_DATE. #1 idempotency bug. - Guard non-replayable side effects (emails, external API mutations).
- Запись с ограничением партицией — пишем только партицию для
{{ ds }}, никогда не «всё с вчерашнего дня». - Перезапись по партиции / MERGE / upsert — перезапуск заменяет, никогда не добавляет. (dbt
--full-refreshна партицию;INSERT OVERWRITE; IcebergreplaceWhere/MERGE.) - Детерминированные входные данные — параметризуй от
logical_date, никогда отdatetime.now()/CURRENT_DATE. Это баг №1 в идемпотентности. - Защищай неповторяемые побочные эффекты (email, мутации внешних API).
# BAD — non-idempotent: appends, uses wall-clock INSERT INTO sales SELECT * FROM src WHERE dt = CURRENT_DATE; # GOOD — idempotent: overwrite the exact partition for the logical date INSERT OVERWRITE sales PARTITION (ds='{{ ds }}') SELECT * FROM src WHERE dt = '{{ ds }}';
# ПЛОХО — неидемпотентно: добавляет, использует реальное время INSERT INTO sales SELECT * FROM src WHERE dt = CURRENT_DATE; # ХОРОШО — идемпотентно: перезаписывает точную партицию для логической даты INSERT OVERWRITE sales PARTITION (ds='{{ ds }}') SELECT * FROM src WHERE dt = '{{ ds }}';
Backfills
| Mechanism | What |
|---|---|
catchup=True | auto-fills missing intervals on unpause |
airflow dags backfill -s -e | bounded historical range on demand |
clear TaskInstances | wipe state → scheduler reruns ("rerun yesterday") |
| Механизм | Что делает |
|---|---|
catchup=True | автоматически заполняет пропущенные интервалы при снятии с паузы |
airflow dags backfill -s -e | ограниченный исторический диапазон по требованию |
clear TaskInstances | стирает состояние → scheduler перезапускает («переиграть вчера») |
All three rely on idempotency. Limit blast radius with max_active_runs + pools so a backfill doesn't DoS sources.
Все три опираются на идемпотентность. Ограничь радиус поражения через max_active_runs + pools, чтобы backfill не задосил источники.
Retries
Ретраи
retries, retry_delay, retry_exponential_backoff=True. Retries help transient failures only if the task is idempotent. Don't retry deterministic failures (bad SQL) — wasted time. Alert on final failure via on_failure_callback.
retries, retry_delay, retry_exponential_backoff=True. Ретраи помогают временным сбоям только если задача идемпотентна. Не ретрай детерминированные ошибки (плохой SQL) — потеря времени. Алерт на окончательный провал через on_failure_callback.
Sensors & XComsSensors и XComs
▾Sensors — wait for a condition
Sensors — ожидание условия
- FileSensor, S3KeySensor, ExternalTaskSensor, custom sensors.
- CRITICAL: use
mode="reschedule"or deferrable operators (triggerer), NOTmode="poke". Poke holds a worker slot the whole time → many sensors = pool starvation / deadlock. - Set
timeout+poke_interval; decide soft-fail vs hard-fail. - Better: prefer event/data-driven triggering (Datasets/Assets) over polling.
- FileSensor, S3KeySensor, ExternalTaskSensor, кастомные sensors.
- КРИТИЧНО: используй
mode="reschedule"или deferrable операторы (triggerer), НЕmode="poke". Poke держит слот воркера всё время → много sensors = истощение pool / deadlock. - Устанавливай
timeout+poke_interval; решай soft-fail vs hard-fail. - Лучше: предпочитай event/data-driven триггеры (Datasets/Assets) вместо поллинга.
XComs — small messages between tasks
XComs — маленькие сообщения между задачами
- Stored in metadata DB; TaskFlow return values auto-become XComs.
- Danger: meant for small metadata (a path, a count, a date). Never push dataframes / large payloads — bloats the DB, kills the scheduler.
- Big data → write to S3/warehouse, pass the pointer via XCom. Custom XCom backend can offload to S3/GCS.
- Хранятся в metadata DB; возвращаемые значения TaskFlow автоматически становятся XComs.
- Опасность: предназначены для малых метаданных (путь, счётчик, дата). Никогда не пиши dataframes / большие payload — раздувает DB, убивает scheduler.
- Большие данные → пиши в S3/хранилище, передавай указатель через XCom. Кастомный XCom-бэкенд может разгружать в S3/GCS.
ExternalTaskSensor — cross-DAG dep trap
ExternalTaskSensor — ловушка cross-DAG зависимостей
Both DAGs must share the same logical date or set execution_delta / execution_date_fn to align. Mismatched schedules ⇒ waits forever. Modern fix: Datasets/Assets decouple from schedule alignment entirely.
Оба DAG должны иметь одинаковую logical date или установить execution_delta / execution_date_fn для выравнивания. Несовпадающие расписания ⇒ ждёт вечно. Современное лечение: Datasets/Assets полностью отвязаны от выравнивания расписаний.
Executors — Scaling the WorkИсполнители — масштабирование работы
▾| Executor | How it runs | Use case | Trade-off |
|---|---|---|---|
| Sequential | 1 task, SQLite | demos | no parallelism |
| Local | subprocesses on scheduler host | small single-node | bounded by one box |
| Celery | queue (Redis/RabbitMQ) + worker fleet | steady horizontal scale | run brokers+workers; idle cost |
| Kubernetes | one pod per task | bursty, isolation, per-task deps | pod startup latency |
| CeleryKubernetes | hybrid | steady + bursty mix | complexity |
| Исполнитель | Как запускает | Применение | Компромисс |
|---|---|---|---|
| Sequential | 1 задача, SQLite | демки | нет параллелизма |
| Local | подпроцессы на хосте scheduler | малый single-node | ограничен одной машиной |
| Celery | очередь (Redis/RabbitMQ) + флот воркеров | стабильная горизонтальная масштаб. | нужны брокеры+воркеры; стоимость простоя |
| Kubernetes | один pod на задачу | пиковая, изоляция, per-task deps | задержка старта pod |
| CeleryKubernetes | гибридный | стабильная + пиковая смесь | сложность |
Concurrency controls: parallelism (cluster) ⊃ max_active_tasks (per DAG) ⊃ Pool slots (per resource) ⊃ priority_weight (queue order).
Контроль параллелизма: parallelism (кластер) ⊃ max_active_tasks (на DAG) ⊃ Pool-слоты (на ресурс) ⊃ priority_weight (порядок в очереди).
Dynamic DAGs, SLAs & TestingДинамические DAG, SLA и тестирование
▾Dynamic — two different things
Dynamic — две разные вещи
- Dynamic DAG generation — build DAG objects in a loop from config (YAML/DB). Cost: parsed every scheduler loop ⇒ keep top-level code cheap (no API/DB calls at module level).
- Dynamic Task Mapping (
.expand(), 2.3+) — fan out a task over a runtime list (map/reduce). Count decided at runtime, visible in UI. Preferred when the count is data-dependent.
- Динамическая генерация DAG — сборка DAG-объектов в цикле из конфига (YAML/DB). Цена: парсится каждый scheduler-цикл ⇒ держи код на верхнем уровне дешёвым (никаких API/DB вызовов на уровне модуля).
- Dynamic Task Mapping (
.expand(), 2.3+) — размножение задачи по runtime-списку (map/reduce). Число решается в runtime, видно в UI. Предпочтительно, когда число зависит от данных.
process = run.expand(item=get_items()) # N instances at runtime# N экземпляров в runtime
SLAs be honest
SLA будь честен
Max expected duration relative to logical date; miss → sla_miss_callback/email. Caveat: built-in SLAs are weak — only evaluated by scheduler, don't fire if the task never starts, tied to logical date. Say: "I'd back it with freshness/deadline checks rather than rely on the built-in feature."
Макс. ожидаемая длительность относительно logical date; промах → sla_miss_callback/email. Оговорка: встроенные SLA слабые — оцениваются только scheduler, не срабатывают, если задача не стартует, привязаны к logical date. Говори: «Я бы подкрепил их проверками freshness/deadline, а не полагался на встроенную фичу».
Testing important
Тестирование важно
- DAG-import tests — assert all DAGs parse, no cycles, expected task counts (pytest over DagBag, run in CI on every PR).
- Unit-test the logic, not the operator — extract business logic into plain functions; operator is glue.
airflow tasks test <dag> <task> <date>— run one TaskInstance without scheduler/DB writes.- Data tests — dbt tests / Great Expectations / row-count + freshness as actual tasks.
- DAG-import тесты — утверждай, что все DAG парсятся, нет циклов, ожидаемое число задач (pytest поверх DagBag, запускай в CI на каждом PR).
- Юнит-тести логику, а не оператор — вынимай бизнес-логику в plain-функции; оператор — клей.
airflow tasks test <dag> <task> <date>— запуск одной TaskInstance без scheduler/DB-записи.- Data-тесты — dbt tests / Great Expectations / row-count + freshness как реальные задачи.
Airflow vs Dataswarm / Chronos vs Spark/FlinkAirflow vs Dataswarm / Chronos vs Spark/Flink
▾Airflow vs Dataswarm/Chronos
Airflow vs Dataswarm/Chronos
| Dimension | Airflow (OSS) | Dataswarm / Chronos (Meta) |
|---|---|---|
| Paradigm | task / control-flow centric | data-asset / partition centric |
| Scheduling | time-interval; you own start_date/catchup | platform-managed; declare partitions + waits |
| Dependencies | explicit task deps + Sensors/Datasets | wait-on-upstream-partition (data-driven) |
| Backfill | dags backfill, you manage it | first-class, platform-orchestrated |
| Infra | you run scheduler/executor/workers/DB | fully managed, autoscaled |
| Observability | UI + wire your own DQ/alerts | integrated lineage, Scuba, DQ monitors |
| Измерение | Airflow (OSS) | Dataswarm / Chronos (Meta) |
|---|---|---|
| Парадигма | task / control-flow центрично | data-asset / partition центрично |
| Расписание | time-interval; ты владеешь start_date/catchup | управляется платформой; декларируешь партиции + waits |
| Зависимости | явные зависимости задач + Sensors/Datasets | wait-on-upstream-partition (data-driven) |
| Backfill | dags backfill, ты управляешь | first-class, оркестрируется платформой |
| Инфраструктура | ты запускаешь scheduler/executor/workers/DB | полностью managed, autoscaled |
| Наблюдаемость | UI + ты подключаешь DQ/alerts | интегрированы lineage, Scuba, DQ monitors |
Where Airflow sits vs compute engines (bridges to DE interview)
Где Airflow по отношению к движкам вычислений (связь с DE-интервью)
Airflow orchestrates batch. Streaming = Flink + Kafka (Airflow is the wrong tool there). Iceberg is the table format under both — its snapshot isolation + MERGE is what makes idempotent partition overwrite clean.
Airflow оркестрирует batch. Стриминг = Flink + Kafka (Airflow там не подходит). Iceberg — формат таблицы под обоими — его snapshot isolation + MERGE делают идемпотентную перезапись партиций чистой.
Top Gotchas / Anti-PatternsОсновные подводные камни / антипаттерны
▾- ▲ Heavy code at top level of DAG files → slow parsing (runs every scheduler loop + every worker).
- ▲ Passing large data through XCom → DB bloat, scheduler death.
- ▲
datetime.now()/CURRENT_DATEinstead of logical date → breaks idempotency & backfills. - ▲ Dynamic
start_date→ DagRuns may never trigger. - ▲
mode="poke"sensors → worker starvation. - ▲ Non-idempotent appends → duplicates on retry.
- ▲ One giant monolithic task vs granular, individually-retryable tasks.
- ▲ Treating task-success as data-correctness (no DQ checks) → silent empty/wrong data.
- ▲ Secrets in Variables → use a secrets backend (Vault / AWS Secrets Manager).
- ▲
catchup=Trueleft on by accident → 100s of concurrent runs hammering sources. - ▲ Renaming a DAG → history keyed on
dag_id; new id loses history + may re-trigger catchup.
- ▲ Тяжёлый код на верхнем уровне DAG-файлов → медленный парсинг (запускается каждый scheduler-цикл + каждый воркер).
- ▲ Передача больших данных через XCom → раздутие DB, смерть scheduler.
- ▲
datetime.now()/CURRENT_DATEвместо logical date → ломает идемпотентность и backfills. - ▲ Динамический
start_date→ DagRuns могут никогда не сработать. - ▲
mode="poke"sensors → истощение воркеров. - ▲ Неидемпотентные appends → дубли при ретраях.
- ▲ Одна гигантская монолитная задача vs гранулярные, по отдельности ретраибл задачи.
- ▲ Считать успех задачи = корректность данных (нет DQ-проверок) → тихие пустые/неправильные данные.
- ▲ Секреты в Variables → используй secrets backend (Vault / AWS Secrets Manager).
- ▲
catchup=Trueоставлен случайно → сотни одновременных запусков молотят источники. - ▲ Переименование DAG → история ключуется на
dag_id; новый id теряет историю + может перезапустить catchup.
Self-Quiz — Flip to RevealСамопроверка — переверни карточку
▾Tap a card to flip. Cover the answer, say it out loud, then check.
Нажми на карточку, чтобы перевернуть. Закрой ответ, произнеси вслух, затем проверь.
@daily run for June 13 fire, and what is its logical_date?@daily за 13 июня, и чему равна его logical_date?now()). So N runs = same end state.now()). Так N запусков = одно конечное состояние.mode="poke" on sensors dangerous at scale?mode="poke" на sensors опасен на масштабе?CURRENT_DATE bug, swallowed exception, schema drift. Success ≠ correctness → add DQ checks.CURRENT_DATE, проглоченное исключение, дрейф схемы. Успех ≠ корректность → добавь DQ-проверки.trigger_rule="all_done" (fires once all upstream finish, any state).trigger_rule="all_done" (срабатывает, как только все upstream завершены, любое состояние)..expand() fans one task over a runtime list (count data-dependent)..expand() размножает одну задачу по runtime-списку (число зависит от данных).start_date?start_date?MERGE/replaceWhere partition overwrite → a rerun just creates a new snapshot replacing that day.MERGE/replaceWhere перезапись партиции → перезапуск просто создаёт новый snapshot, заменяя тот день.