Airflow & OrchestrationAirflow и оркестрация

Fast-revision cheatsheet for Data Engineering interviews. Tailored for Senior/Lead DE. Шпаргалка для быстрого повторения перед Data Engineering интервью. Заточена под Senior/Lead DE.
⭐ TOP topic: idempotency + backfills⭐ ТОПОВАЯ тема: идемпотентность + backfills Anchor story: McMakler Airflow migrationЯкорная история: миграция на Airflow в McMakler Study-gap bridge: Flink / Kafka / IcebergПробел для изучения: Flink / Kafka / Iceberg
Your one-liner credential At McMakler you migrated a proprietary orchestrator to Airflow, killed 100s of K EUR/yr in licensing, built a proprietary Salesforce sync layer on Airflow, and grew the team 2→7. That is "I have actually operated Airflow in production." Lead with it.
Твоя формула в одно предложение В McMakler ты мигрировал проприетарный оркестратор на Airflow, убил сотни тысяч EUR/год на лицензиях, построил проприетарный слой синхронизации Salesforce на Airflow и вырастил команду с 2 до 7 человек. Это «я реально эксплуатировал Airflow в продакшене». Начинай с этого.

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/контейнеры, которые это делают.

WHAT AIRFLOW IS vs WHAT IT IS NOT ───────────────────────── ───────────────────────── • schedules work (time/data) ✗ a data processing engine • declares dependencies (DAG) ✗ a streaming system (use Flink/Kafka) • retries + tracks task state ✗ a guarantee of correct data • triggers Spark/dbt/SQL/pods ✗ exactly-once (it's at-least-once)
ЧТО ТАКОЕ AIRFLOW vs ЧТО ОН НЕ ДЕЛАЕТ ───────────────────────── ───────────────────────── • планирует работу (время/данные) ✗ движок обработки данных • объявляет зависимости (DAG) ✗ стриминг-система (используй Flink/Kafka) • ретраи + отслеживание состояния ✗ гарантия корректных данных • запускает Spark/dbt/SQL/pods ✗ exactly-once (это at-least-once)

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Жизненный цикл выполнения

DAG file (.py) │ parsed by ▼ SCHEDULER / DAG Processor ──► creates DagRun when data_interval elapses │ creates TaskInstance (scheduled→queued) ▼ EXECUTOR (Local / Celery / Kubernetes) │ dispatches to ▼ WORKER ──► re-parses DAG file, runs operator.execute() in fresh process │ ▼ METADATA DB (Postgres) ◄── state: running→success/failed/up_for_retry ▲ scheduler polls to release downstream tasks │ WEBSERVER / UI (reads metadata DB)
DAG-файл (.py) │ парсится ▼ SCHEDULER / DAG Processor ──► создаёт DagRun по окончании data_interval │ создаёт TaskInstance (scheduled→queued) ▼ EXECUTOR (Local / Celery / Kubernetes) │ отправляет на ▼ WORKER ──► заново парсит DAG-файл, запускает operator.execute() в новом процессе │ ▼ METADATA DB (Postgres) ◄── состояние: running→success/failed/up_for_retry ▲ scheduler опрашивает, чтобы освободить downstream-таски │ WEBSERVER / UI (читает metadata DB)
Interviewer will probe "When does a @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Базовые объекты и терминология

TermWhat it is
DAGDirected Acyclic Graph: the pipeline definition. Acyclic ⇒ valid topological order, no infinite loops.
DagRunOne execution of a DAG for a specific data interval / logical date.
OperatorTemplate/class for a unit of work (PythonOperator, BashOperator, SQL ops, KubernetesPodOperator).
TaskAn operator instantiated in a DAG (a node in the graph).
TaskInstanceOne run of one task for one DagRun — the thing that actually has state.
SensorSpecial operator that waits for a condition (file, partition, external task).
XComCross-communication: small key/value passed between tasks via metadata DB.
Connection / VariableCredentials+endpoints / arbitrary global config k-v.
PoolCaps concurrent task slots for a scarce resource (e.g. a fragile DB).
TriggererAsync process running deferrable operators so waiting tasks don't hold worker slots.
ТерминЧто это
DAGDirected Acyclic Graph: определение пайплайна. Acyclic ⇒ валидный топологический порядок, нет бесконечных циклов.
DagRunОдно выполнение DAG для конкретного data interval / logical date.
OperatorШаблон/класс для единицы работы (PythonOperator, BashOperator, SQL ops, KubernetesPodOperator).
TaskОператор, проинстанцированный в DAG (узел в графе).
TaskInstanceОдин запуск одной задачи для одного DagRun — вот он-то и имеет состояние.
SensorСпециальный оператор, который ждёт условия (файл, партиция, внешняя задача).
XComCross-communication: маленький key/value, передаваемый между задачами через metadata DB.
Connection / VariableCredentials+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.
  • catchupTrue backfills every interval from start_date on unpause (can spawn 100s of runs). Default to catchup=False + controlled backfills.
  • max_active_runs / max_active_tasks / parallelism — concurrency blast-radius control.
  • start_date — первый рассматриваемый интервал. Никогда динамический (datetime.now()) → неопределённый первый интервал.
  • schedule — cron / пресет (@daily) / timedelta / None / список dataset.
  • catchupTrue заполняет каждый интервал от 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 ruleFires when…Use for
all_success (default)all upstream succeedednormal flow
all_doneall upstream finished (any state)cleanup / notify regardless
none_failed_min_one_successnone failed, ≥1 succeededjoin after branching
one_failed / one_success≥1 in that statefast-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 алерты
Interviewer will probe "Run a notify task even if the pipeline failed?" → trigger_rule="all_done". "Conditional path?" → @task.branch + none_failed_min_one_success on the join.
О чём спросит интервьюер «Запустить notify-таску, даже если пайплайн упал?» → 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-refresh on a partition; INSERT OVERWRITE; Iceberg replaceWhere / MERGE.)
  • Deterministic inputs — parameterize off logical_date, never datetime.now() / CURRENT_DATE. #1 idempotency bug.
  • Guard non-replayable side effects (emails, external API mutations).
  • Запись с ограничением партицией — пишем только партицию для {{ ds }}, никогда не «всё с вчерашнего дня».
  • Перезапись по партиции / MERGE / upsert — перезапуск заменяет, никогда не добавляет. (dbt --full-refresh на партицию; INSERT OVERWRITE; Iceberg replaceWhere / 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

MechanismWhat
catchup=Trueauto-fills missing intervals on unpause
airflow dags backfill -s -ebounded historical range on demand
clear TaskInstanceswipe 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.

Anchor to your experience At Meta you fixed a green-but-wrong defect across 6.2B records / 34% of threads and restored 100+ downstream metrics. Use it for "success ≠ correctness." At McMakler, the Salesforce sync had to be idempotent — re-syncing must upsert on a natural key, never duplicate records.
Привязка к опыту В Meta ты исправил дефект «зелёно, но неправильно» на 6,2 млрд записей / 34% тредов и восстановил 100+ downstream-метрик. Используй это для «успех ≠ корректность». В McMakler синхронизация Salesforce должна была быть идемпотентной — ре-синк должен делать upsert по натуральному ключу, никогда не дублировать записи.

Sensors & XComsSensors и XComs

Sensors — wait for a condition

Sensors — ожидание условия

  • FileSensor, S3KeySensor, ExternalTaskSensor, custom sensors.
  • CRITICAL: use mode="reschedule" or deferrable operators (triggerer), NOT mode="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) вместо поллинга.
Interviewer will probe "10,000 sensors waiting on files — what breaks?" → poke-mode worker starvation. Fix: reschedule mode / deferrable operators that offload to the async triggerer (near-zero resources while waiting).
О чём спросит интервьюер «10 000 sensors ждут файлов — что сломается?» → poke-режим истощает воркеры. Лечение: reschedule-режим / deferrable операторы, которые разгружают на асинхронный triggerer (почти нулевые ресурсы при ожидании).

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Исполнители — масштабирование работы

ExecutorHow it runsUse caseTrade-off
Sequential1 task, SQLitedemosno parallelism
Localsubprocesses on scheduler hostsmall single-nodebounded by one box
Celeryqueue (Redis/RabbitMQ) + worker fleetsteady horizontal scalerun brokers+workers; idle cost
Kubernetesone pod per taskbursty, isolation, per-task depspod startup latency
CeleryKuberneteshybridsteady + bursty mixcomplexity
ИсполнительКак запускаетПрименениеКомпромисс
Sequential1 задача, SQLiteдемкинет параллелизма
Localподпроцессы на хосте schedulerмалый single-nodeограничен одной машиной
Celeryочередь (Redis/RabbitMQ) + флот воркеровстабильная горизонтальная масштаб.нужны брокеры+воркеры; стоимость простоя
Kubernetesодин pod на задачупиковая, изоляция, per-task depsзадержка старта pod
CeleryKubernetesгибридныйстабильная + пиковая смесьсложность
Interviewer will probe "Workers idle off-peak, cost matters" → Kubernetes (pod-per-task, scales to zero) or deferrable operators. "High throughput, steady, low latency" → Celery (warm workers, no pod cold-start).
О чём спросит интервьюер «Воркеры простаивают в низкую нагрузку, важна цена» → Kubernetes (pod-per-task, масштабируется до нуля) или deferrable операторы. «Высокая пропускная, стабильная, низкая задержка» → Celery (тёплые воркеры, нет холодного старта pod).

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

DimensionAirflow (OSS)Dataswarm / Chronos (Meta)
Paradigmtask / control-flow centricdata-asset / partition centric
Schedulingtime-interval; you own start_date/catchupplatform-managed; declare partitions + waits
Dependenciesexplicit task deps + Sensors/Datasetswait-on-upstream-partition (data-driven)
Backfilldags backfill, you manage itfirst-class, platform-orchestrated
Infrayou run scheduler/executor/workers/DBfully managed, autoscaled
ObservabilityUI + wire your own DQ/alertsintegrated lineage, Scuba, DQ monitors
ИзмерениеAirflow (OSS)Dataswarm / Chronos (Meta)
Парадигмаtask / control-flow центричноdata-asset / partition центрично
Расписаниеtime-interval; ты владеешь start_date/catchupуправляется платформой; декларируешь партиции + waits
Зависимостиявные зависимости задач + Sensors/Datasetswait-on-upstream-partition (data-driven)
Backfilldags backfill, ты управляешьfirst-class, оркестрируется платформой
Инфраструктураты запускаешь scheduler/executor/workers/DBполностью managed, autoscaled
НаблюдаемостьUI + ты подключаешь DQ/alertsинтегрированы lineage, Scuba, DQ monitors
Say this "Dataswarm pushes a declarative, data-asset model — does the output partition exist? — whereas vanilla Airflow is control-flow centric. But Airflow's Datasets/Assets are converging on the same model. I've seen both: managed (Meta) and operate-it-yourself (McMakler). The principles — idempotency, partitioning, dependency-as-data, backfills — transfer directly."
Скажи это «Dataswarm продвигает декларативную, data-asset модель — существует ли выходная партиция? — тогда как обычный Airflow центрируется на control-flow. Но Airflow-овые Datasets/Assets сходятся к той же модели. Я видел оба: managed (Meta) и operate-it-yourself (McMakler). Принципы — идемпотентность, партиционирование, dependency-as-data, backfills — переносятся напрямую.»

Where Airflow sits vs compute engines (bridges to DE interview)

Где Airflow по отношению к движкам вычислений (связь с DE-интервью)

Kafka ──► Flink (stream, stateful, exactly-once via checkpoints) ──► IcebergAirflow ──► Spark / dbt / SQL (batch / ELT) ──────────────────────┘ orchestrates batch compute open table format (ACID, snapshots, time-travel)

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_DATE instead 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=True left 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.
Study-gap — be clear & honest For Flink / Iceberg: "My production Airflow/Spark/Kafka experience is real (McMakler migration, Front Tier Scala/Spark-on-Hadoop, IU Group Kafka). Flink & Iceberg I know conceptually and am ramping." Then connect: Iceberg snapshot isolation = clean idempotent overwrites; Flink = streaming counterpart of batch Spark (event-time, watermarks, checkpoints for exactly-once); Kafka = the log Airflow batch-consumes or triggers jobs against.
Пробел для изучения — будь чёток и честен По Flink / Iceberg: «Мой продакшен-опыт Airflow/Spark/Kafka реален (миграция McMakler, Front Tier Scala/Spark-на-Hadoop, IU Group Kafka). Flink и Iceberg я знаю концептуально и разгоняюсь.» Затем свяжи: Iceberg snapshot isolation = чистая идемпотентная перезапись; Flink = стриминговый аналог batch Spark (event-time, watermarks, checkpoints для exactly-once); Kafka = лог, который Airflow batch-консумит или против которого запускает jobs.

Self-Quiz — Flip to RevealСамопроверка — переверни карточку

Tap a card to flip. Cover the answer, say it out loud, then check.

Нажми на карточку, чтобы перевернуть. Закрой ответ, произнеси вслух, затем проверь.

When does a @daily run for June 13 fire, and what is its logical_date?
Когда запускается @daily за 13 июня, и чему равна его logical_date?
tap to flip
нажми, чтобы перевернуть
Fires at the END of the interval (00:00 June 14). logical_date = start of interval = June 13. Lags wall-clock by one interval.
Срабатывает в КОНЦЕ интервала (00:00 14 июня). logical_date = начало интервала = 13 июня. Отстаёт от реального времени на один интервал.
How do you make a task idempotent?
Как сделать задачу идемпотентной?
tap to flip
нажми, чтобы перевернуть
Partition-scoped writes keyed on logical_date + overwrite/MERGE/upsert (not append) + deterministic inputs (no now()). So N runs = same end state.
Запись по партициям на logical_date + overwrite/MERGE/upsert (не append) + детерминированные входы (не now()). Так N запусков = одно конечное состояние.
Why is mode="poke" on sensors dangerous at scale?
Почему mode="poke" на sensors опасен на масштабе?
tap to flip
нажми, чтобы перевернуть
Each sensor holds a worker slot the entire wait → pool starvation/deadlock. Use reschedule mode or deferrable operators (async triggerer).
Каждый sensor держит слот воркера всё время ожидания → истощение pool/deadlock. Используй reschedule mode или deferrable операторы (async triggerer).
What must NOT go through XCom, and why?
Что НЕЛЬЗЯ передавать через XCom, и почему?
tap to flip
нажми, чтобы перевернуть
Large payloads / dataframes — XComs live in the metadata DB. Pass a pointer (S3 path); store the data in the warehouse.
Большие payload / dataframes — XComs хранятся в metadata DB. Передавай указатель (путь S3); храни данные в хранилище.
Daily DAG is green for 3 days but produced no data. First checks?
Ежедневный DAG зелёный 3 дня, но не создал данных. Первые проверки?
tap to flip
нажми, чтобы перевернуть
Are DagRuns even created? (paused / scheduler / catchup). If tasks "succeeded" empty → upstream empty, CURRENT_DATE bug, swallowed exception, schema drift. Success ≠ correctness → add DQ checks.
Вообще создаются DagRuns? (paused / scheduler / catchup). Если задачи «успешны» пусто → upstream пуст, баг CURRENT_DATE, проглоченное исключение, дрейф схемы. Успех ≠ корректность → добавь DQ-проверки.
Executor for bursty load where workers must scale to zero off-peak?
Исполнитель для пиковой нагрузки, где воркеры должны масштабироваться до нуля в низкую нагрузку?
tap to flip
нажми, чтобы перевернуть
KubernetesExecutor (pod per task, scales to zero) or deferrable operators. Celery = steady high-throughput, low latency.
KubernetesExecutor (pod на задачу, масштабируется до нуля) или deferrable операторы. Celery = стабильная высокая пропускная, низкая задержка.
Run a notify/cleanup task even if upstream failed?
Запустить notify/cleanup-задачу, даже если upstream упал?
tap to flip
нажми, чтобы перевернуть
trigger_rule="all_done" (fires once all upstream finish, any state).
trigger_rule="all_done" (срабатывает, как только все upstream завершены, любое состояние).
Consistency guarantee Airflow gives?
Какую гарантию консистентности даёт Airflow?
tap to flip
нажми, чтобы перевернуть
At-least-once execution (retries/reruns/backfills). Exactly-once outputs only via idempotent tasks.
At-least-once execution (ретраи/перезапуски/backfills). Exactly-once выходы только через идемпотентные задачи.
Dynamic DAG generation vs Dynamic Task Mapping?
Динамическая генерация DAG vs динамическое маппинг задач?
tap to flip
нажми, чтобы перевернуть
Generation = build N DAG objects from config (cost: parsed every loop). Mapping = .expand() fans one task over a runtime list (count data-dependent).
Генерация = сборка N DAG-объектов из конфига (цена: парсится каждый цикл). Маппинг = .expand() размножает одну задачу по runtime-списку (число зависит от данных).
Airflow vs Dataswarm — one-sentence framing?
Airflow vs Dataswarm — формула в одно предложение?
tap to flip
нажми, чтобы перевернуть
Airflow = control-flow/time centric (you operate it); Dataswarm = managed, data-asset/partition centric. Airflow's Datasets are converging on that model; principles transfer.
Airflow = control-flow/время центричен (ты эксплуатируешь); Dataswarm = managed, data-asset/partition центричен. Airflow-овые Datasets сходятся к той модели; принципы переносятся.
Why never use a dynamic start_date?
Почему никогда не использовать динамический start_date?
tap to flip
нажми, чтобы перевернуть
It makes the first data interval undefined; the scheduler may never compute a triggerable interval → DagRuns don't fire.
Это делает первый data interval неопределённым; scheduler может никогда не вычислить запускаемый интервал → DagRuns не срабатывают.
How does Iceberg make idempotent writes clean?
Как Iceberg делает идемпотентные записи чистыми?
tap to flip
нажми, чтобы перевернуть
Snapshot isolation (readers never see partial partition) + MERGE/replaceWhere partition overwrite → a rerun just creates a new snapshot replacing that day.
Snapshot isolation (читатели не видят частичную партицию) + MERGE/replaceWhere перезапись партиции → перезапуск просто создаёт новый snapshot, заменяя тот день.