Оркестрация пайплайнов: Airflow, dbt, идемпотентность и backfill
К этому моменту трека у нас есть все кирпичи: мы умеем забирать данные (https://courses.digitable.life/post/data-engineering/01-etl-vs-elt/), раскладывать их в модель (https://courses.digitable.life/post/data-engineering/02-data-modeling/), считать батчем (https://courses.digitable.life/post/data-engineering/03-batch-processing/) и стримом (https://courses.digitable.life/post/data-engineering/04-streaming/). Но кирпичи — ещё не дом. В реальной компании таких вычислений сотни, они зависят друг от друга, падают, требуют пересчёта за прошлое, конкурируют за одну и ту же реплику базы и должны быть готовы к 9:00, когда директор откроет дашборд.
Оркестрация — это дисциплина, которая отвечает на вопрос: что и когда запускать, что делать при падении и как безопасно пересчитать прошлое. Оркестратор не считает данные сам; он решает, кому дать команду считать, в каком порядке и при каких условиях.
Аналогия: повар готовит блюдо, но в ресторане есть ещё и экспедитор — человек, который смотрит на все заказы сразу, знает, что гарнир должен быть готов одновременно с мясом, что плита занята и что если соус пригорел — переделать надо соус, а не всё блюдо. Оркестратор — экспедитор кухни данных.
1. Модель: DAG как язык зависимостей
Базовая абстракция оркестрации — направленный ациклический граф (DAG, directed acyclic graph). Вершины — единицы работы (задачи), рёбра — отношение «должно завершиться до».
Почему DAG, а не список шагов? Список слишком строг: если загрузить_заказы и загрузить_клиентов
независимы, линейный порядок заставит их ждать друг друга, а DAG выражает параллелизм бесплатно.
Ацикличность же — не каприз, а условие существования топологической сортировки: если A ждёт B, а B
ждёт A, плана не существует.
Планировщик оркестратора — это, по сути, топологическая сортировка с ограничением ресурсов. Базовый алгоритм (алгоритм Кана) даёт корректный порядок за O(V + E):
from collections import deque
def kahn_schedule(tasks: dict[str, list[str]]) -> list[str]:
"""tasks: задача -> список её зависимостей (upstream).
Корректный порядок выполнения за O(V + E) по времени и памяти."""
downstream = {t: [] for t in tasks}
indegree = {t: 0 for t in tasks}
for task, deps in tasks.items():
for dep in deps:
downstream[dep].append(task) # ребро dep -> task
indegree[task] += 1
ready = deque(t for t, d in indegree.items() if d == 0) # без невыполненных зависимостей
order = []
while ready:
task = ready.popleft()
order.append(task)
for nxt in downstream[task]:
indegree[nxt] -= 1
if indegree[nxt] == 0:
ready.append(nxt)
if len(order) != len(tasks):
raise ValueError("В графе есть цикл — план невыполним")
return order
Реальный планировщик усложняет это в трёх местах, и именно там живёт вся инженерия:
- Ресурсы. Готовых задач может быть 200, а слотов — 20. «Выбрать оптимальный порядок при ограниченных ресурсах» — это job-shop scheduling, NP-трудная задача, поэтому все оркестраторы используют жадные эвристики (приоритет, FIFO внутри приоритета) и пулы.
- Отказы. Задача не просто «выполнена/не выполнена» — есть ретраи, таймауты, зомби-процессы.
- Время. Граф запускается не один раз, а по расписанию: получается сетка «граф × момент времени». Именно из-за этого измерения появляются backfill и идемпотентность.
Вот тот же пайплайн как граф — обратите внимание на естественный параллелизм двух веток:
Postgres → S3"] EC["extract_customers
CRM API → S3"] end subgraph TRANSFORM["Transform — dbt"] SO["stg_orders"] SC["stg_customers"] DC["dim_customer
SCD2"] FO["fct_orders"] end subgraph SERVE["Публикация"] Q["dq_checks
тесты качества"] P["publish_bi
обновить кэш BI"] end EO --> SO EC --> SC SC --> DC SO --> FO DC --> FO FO --> Q Q -->|"все тесты зелёные"| P Q -.->|"провал"| AL["alert в дежурный канал"]
Важная деталь: dq_checks стоит между расчётом и публикацией. Это шаблон «ворота качества» —
плохие данные не должны доехать до потребителя. Подробно о контрактах и проверках — в статье
https://courses.digitable.life/post/data-engineering/07-data-quality-and-governance/.
2. Время: логическая дата против времени запуска
Это место, где ломается больше всего пайплайнов. Расписание @daily обрабатывает данные за сутки,
но сутки заканчиваются в полночь: запуск, считающий 11 февраля, физически стартует 12-го. Возникают
четыре разных момента времени, и путать их нельзя:
| Понятие | Что означает | Пример |
|---|---|---|
data_interval_start |
начало окна данных | 2026-02-11 00:00 |
data_interval_end |
конец окна (исключая) | 2026-02-12 00:00 |
logical_date |
метка запуска (в Airflow = начало интервала) | 2026-02-11 |
| фактический старт | когда планировщик реально дёрнул задачу | 2026-02-12 00:07 |
Золотое правило: задача не должна знать, «который сейчас час». Любое обращение к now(),
CURRENT_DATE, datetime.today() внутри логики пайплайна — баг. Причины:
- при backfill за прошлый месяц
now()вернёт сегодня, и вы тридцать раз перезапишете одну и ту же партицию сегодняшним днём; - при ретрае в 00:59 и в 01:03 задача может попасть в разные сутки;
- результат перестаёт быть воспроизводимым — вы не можете доказать, что вчерашняя цифра была верной.
Вместо этого время инжектится снаружи как параметр запуска — задача принимает
data_interval_start / data_interval_end и больше нигде не смотрит на часы.
И ещё: интервал обязан быть полуоткрытым, [start, end). Если написать BETWEEN start AND end,
событие ровно в 00:00:00 попадёт и во вчерашний, и в сегодняшний запуск. Дубликат ровно один раз в
сутки — самый мучительный тип бага, потому что он маленький и стабильный.
3. Идемпотентность: главное свойство пайплайна
Идемпотентность — свойство операции, при котором повторное применение с теми же аргументами не
меняет результат: f(f(x)) = f(x). Для пайплайна это значит: запустив задачу за 11 февраля один раз,
пять раз или пятьдесят, вы получите в целевой таблице одно и то же состояние.
Зачем это настолько важно? Потому что перезапуск неизбежен: упала сеть, кончилось место, воркер убили при деплое, нашли баг в SQL, аналитик попросил пересчитать квартал. Без идемпотентности каждый такой случай превращается в ручную операцию «сначала удалите вот это, потом запустите вот то» — а ручные операции в 3 часа ночи заканчиваются испорченным хранилищем.
Формально нужна детерминированная область записи, зависящая только от параметров запуска: тогда перезапись этой области полностью стирает следы предыдущей попытки.
Три рабочих паттерна
1. Overwrite партиции (лучший вариант для батча). Задача за 11 февраля владеет ровно партицией
dt=2026-02-11 и переписывает её целиком:
-- Spark / Hive / Trino: динамическая перезапись только затронутых партиций
-- ВАЖНО: spark.sql.sources.partitionOverwriteMode=dynamic,
-- иначе INSERT OVERWRITE снесёт ВСЮ таблицу, а не одну партицию.
INSERT OVERWRITE TABLE analytics.fct_orders PARTITION (dt)
SELECT o.order_id, o.customer_id, o.amount, date(o.created_at) AS dt
FROM staging.orders o
WHERE o.created_at >= TIMESTAMP '{{ data_interval_start }}'
AND o.created_at < TIMESTAMP '{{ data_interval_end }}';
Сложность: чтение источника за окно + запись окна, то есть O(размер окна), а не O(размера таблицы) — поэтому пересчёт одного дня стоит одинаково и в первый день жизни таблицы, и на третий год.
2. MERGE / UPSERT по бизнес-ключу. Когда данные обновляются задним числом и партиция не совпадает с окном загрузки:
MERGE INTO analytics.dim_customer AS t
USING (
SELECT customer_id, name, segment, updated_at
FROM staging.customers
WHERE updated_at >= TIMESTAMP '{{ data_interval_start }}' - INTERVAL '1' HOUR -- нахлёст
AND updated_at < TIMESTAMP '{{ data_interval_end }}'
QUALIFY row_number() OVER (PARTITION BY customer_id ORDER BY updated_at DESC) = 1
) AS s
ON t.customer_id = s.customer_id
WHEN MATCHED AND s.updated_at > t.updated_at THEN UPDATE SET
name = s.name, segment = s.segment, updated_at = s.updated_at
WHEN NOT MATCHED THEN INSERT (customer_id, name, segment, updated_at)
VALUES (s.customer_id, s.name, s.segment, s.updated_at);
Два нюанса, без которых MERGE ломается: дедупликация источника (QUALIFY row_number() = 1) — иначе
многие движки бросят ошибку «multiple source rows matched», — и условие s.updated_at > t.updated_at,
чтобы повторный прогон старых данных не откатил таблицу назад.
3. Ключ дедупликации на приёмнике. Для стриминга и внешних API, где перезапись невозможна:
генерируем детерминированный идентификатор из содержимого (sha256(source, business_key, interval))
и вставляем с ON CONFLICT DO NOTHING. Это превращает «at least once» доставку в эффективно
«exactly once» на уровне состояния — см. разбор семантик в https://courses.digitable.life/post/data-engineering/04-streaming/.
Главный антипаттерн — голый INSERT ... SELECT в целевую таблицу без предварительного удаления окна.
Это источник дублей №1 в индустрии, и симптом узнаваем: метрика внезапно ровно вдвое больше нормы за
один конкретный день.
Наконец, идемпотентность записи в таблицу — только половина дела. Задача может ещё отправить письмо,
дёрнуть вебхук или списать деньги. Такие эффекты выносят в отдельную задачу в конце DAG и защищают
ключом идемпотентности на стороне приёмника (Idempotency-Key в HTTP — стандарт платёжных API,
Stripe: Idempotent requests).
4. Жизненный цикл задачи
Задача — не «выполнена / не выполнена». В любом зрелом оркестраторе это конечный автомат, и понимание его состояний экономит часы отладки:
Три состояния тут неочевидны:
queued≠running. Задача может часами висеть вqueued, если пул забит или воркеров мало. Разница между «пайплайн медленный» и «пайплайн голодает по слотам» видна только здесь.zombie. Процесс воркера убит (OOM, спот-инстанс отобрали), но запись в метабазе всё ещёrunning. Планировщик замечает пропажу heartbeat и возвращает задачу в ретрай — и если задача не идемпотентна, это восстановление создаст дубли.up_for_reschedule. Сенсор в режимеpokeдержит слот всё время ожидания, вreschedule— освобождает между проверками. Классический «дедлок»: 32 poke-сенсора заняли все 32 слота и ждут задач, которые некому запустить.
5. Airflow: как это выглядит на практике
Apache Airflow — де-факто стандарт оркестрации данных с 2015 года (проект создан в Airbnb, автор — Maxime Beauchemin). Его ключевая идея — DAG как код на Python, а не как XML/GUI-конфиг: граф можно генерировать циклом, покрывать тестами, ревьюить в PR.
Компонентов пять: DAG Processor парсит .py-файлы, Metadata DB (Postgres) хранит состояние,
Scheduler решает, что запускать, Executor + воркеры выполняют, Webserver показывает.
Отсюда два следствия, которые определяют почти всю эксплуатацию.
Worker и scheduler общаются только через метабазу. Она — единая точка правды и одновременно узкое
место: на тысячах задач в минуту Postgres оркестратора становится главным ограничителем, и его тюнинг
(пулы соединений, autovacuum, чистка старых task_instance) — постоянная работа платформенной команды.
DAG-файл парсится регулярно, каждые несколько десятков секунд. Запрос к базе на верхнем уровне файла выполняется сотни раз в час, парсинг упирается в таймаут, и DAG «пропадает» из UI. Верхний уровень должен быть дешёвым — только описание структуры.
Рабочий DAG
"""Ежедневная витрина заказов. Airflow 2.x, TaskFlow API."""
from __future__ import annotations
import pendulum
from airflow.decorators import dag, task
from airflow.providers.common.sql.operators.sql import SQLExecuteQueryOperator
from airflow.datasets import Dataset
# Датасет — декларативная зависимость между DAG'ами:
# потребитель запустится сам, когда продюсер обновит этот датасет.
FCT_ORDERS = Dataset("warehouse://analytics/fct_orders")
DEFAULT_ARGS = {
"owner": "data-platform",
"retries": 3,
# экспоненциальная выдержка: 5м, 10м, 20м — не долбим упавший источник
"retry_delay": pendulum.duration(minutes=5),
"retry_exponential_backoff": True,
"max_retry_delay": pendulum.duration(hours=1),
"execution_timeout": pendulum.duration(hours=2), # обязательно: иначе задача висит вечно
}
@dag(
dag_id="orders_daily",
schedule="0 3 * * *", # 03:00 — даём источнику время долить данные
start_date=pendulum.datetime(2026, 1, 1, tz="Europe/Moscow"),
catchup=True, # прошлые интервалы будут досчитаны автоматически
max_active_runs=1, # запуски не наезжают друг на друга
default_args=DEFAULT_ARGS,
tags=["orders", "daily", "core"],
doc_md=__doc__,
)
def orders_daily():
@task(pool="postgres_replica", pool_slots=1) # пул ограничивает нагрузку на реплику
def extract_orders(data_interval_start=None, data_interval_end=None) -> str:
"""Выгружает окно [start, end) в S3. Путь детерминирован → перезапуск перезапишет тот же файл."""
import pandas as pd
from sqlalchemy import create_engine
path = f"s3://raw/orders/dt={data_interval_start.date()}/data.parquet"
sql = ("SELECT order_id, customer_id, amount, created_at FROM orders "
"WHERE updated_at >= %(start)s AND updated_at < %(end)s")
df = pd.read_sql(sql, create_engine("postgresql://reader@replica/prod"),
params={"start": data_interval_start, "end": data_interval_end})
if df.empty: # ворота качества: лучше упасть, чем дыра в витрине
raise ValueError(f"Пустая выгрузка {path}: источник вероятно недоступен")
df.to_parquet(path, index=False) # перезапись, не append
return path
load = SQLExecuteQueryOperator(
task_id="load_fct_orders",
conn_id="warehouse",
# Идемпотентно: сначала удаляем окно, потом вставляем. В одной транзакции.
sql="""
BEGIN;
DELETE FROM analytics.fct_orders WHERE dt = DATE '{{ data_interval_start | ds }}';
INSERT INTO analytics.fct_orders (order_id, customer_id, amount, dt)
SELECT order_id, customer_id, amount, DATE '{{ data_interval_start | ds }}'
FROM external.raw_orders
WHERE dt = DATE '{{ data_interval_start | ds }}';
COMMIT;
""",
outlets=[FCT_ORDERS], # сигнал downstream-DAG'ам
)
extract_orders() >> load
orders_daily()
Что здесь стоит скопировать в свои DAG’и:
execution_timeoutвсегда. Задача без таймаута однажды зависнет на сокете и будет держать слот неделю.poolдля всего, что ходит во внешнюю систему. Пул — семафор: «не более 5 одновременных подключений к реплике». Без него backfill за год откроет 365 соединений и положит прод.max_active_runs=1для пайплайнов с состоянием (SCD2, накопительные витрины). Иначе два запуска одновременно перепишут одну область.- Экспоненциальный backoff. Ретраи с фиксированной задержкой в 10 секунд — это DDoS собственного API в момент, когда ему и так плохо.
- XCom только для метаданных. Между задачами передаётся путь к файлу, а не сам DataFrame: XCom хранится в метабазе, и попытка положить туда 200 МБ убьёт Postgres оркестратора.
Сенсоры и деферред-операторы
Ожидание внешнего события («файл появился», «партиция в источнике готова») — отдельная тема. Наивный сенсор занимает слот воркера всё время ожидания. Решения по возрастанию зрелости:
# 1. Плохо: держит слот воркера все 6 часов ожидания
FileSensor(task_id="wait", filepath="/data/flag", poke_interval=60)
# 2. Лучше: mode="reschedule" — между проверками задача уходит из воркера
FileSensor(task_id="wait", filepath="/data/flag", poke_interval=300,
mode="reschedule", timeout=6 * 3600)
# 3. Ещё лучше: deferrable — ожидание живёт в triggerer на asyncio,
# тысячи ожиданий стоят как один процесс
S3KeySensor(task_id="wait", bucket_key="s3://raw/orders/_SUCCESS",
deferrable=True, timeout=6 * 3600)
# 4. Лучше всего: не опрашивать вообще — событийная зависимость
@dag(schedule=[FCT_ORDERS]) # запустится, когда продюсер обновит датасет
def downstream_dag(): ...
Правило: push вместо pull. Опрос — компромисс для систем, которые не умеют уведомлять.
6. dbt: оркестрация внутри трансформаций
Airflow оркестрирует задачи. Но внутри хранилища сотни SQL-моделей со своими зависимостями, и
описывать каждую отдельной задачей Airflow — путь к DAG на 800 вершин, который никто не понимает.
dbt решает это иначе: вы пишете SELECT, а зависимости выводятся
автоматически из ссылок ref() — dbt строит внутренний DAG моделей и оборачивает SELECT в DDL.
-- models/marts/fct_orders.sql
{{ config(
materialized='incremental',
incremental_strategy='delete+insert', -- удалить окно и вставить заново = идемпотентно
unique_key='order_id',
partition_by={'field': 'order_date', 'data_type': 'date'}
) }}
select
o.order_id, o.customer_id, c.segment, o.amount,
date(o.created_at) as order_date
from {{ ref('stg_orders') }} o -- ref() = ребро графа, порядок dbt выведет сам
left join {{ ref('dim_customer') }} c on o.customer_id = c.customer_id
{% if is_incremental() %}
-- только окно запуска. run_date передаёт Airflow — НЕ current_date, иначе backfill сломается
where date(o.created_at) = date '{{ var("run_date") }}'
{% endif %}
Проверки качества объявляются декларативно рядом с моделью:
# models/marts/schema.yml
version: 2
models:
- name: fct_orders
description: "Витрина заказов, гранулярность — один заказ"
columns:
- name: order_id
tests: [unique, not_null]
- name: customer_id
tests:
- not_null
- relationships: # ссылочная целостность между моделями
to: ref('dim_customer')
field: customer_id
Как склеить dbt и Airflow
Три уровня гранулярности — выбор влияет на отладку сильнее, чем кажется:
| Подход | Что в DAG Airflow | Плюсы | Минусы |
|---|---|---|---|
| Монозадача | один dbt build |
просто, быстро внедрить | падение любой модели = красный квадрат без деталей; ретрай пересчитывает всё |
| По группам | задача на слой (staging/marts) или на тег | компромисс, обычно достаточно | грубая зернистость ретрая |
| Модель = задача | генерация из manifest.json (Cosmos, dbt ls) |
точный ретрай, наблюдаемость на уровне модели | DAG на сотни вершин, нагрузка на планировщик |
Практичный рецепт для большинства команд — средний вариант, с явными «воротами»:
def dbt(task_id: str, cmd: str) -> BashOperator:
# {{ ds }} — логическая дата Airflow, прокинутая в переменную dbt
return BashOperator(task_id=task_id,
bash_command=f"dbt {cmd} --vars '{{\"run_date\": \"{{{{ ds }}}}\"}}'")
dbt("run_staging", "run --select tag:staging") >> \
dbt("test_staging", "test --select tag:staging") >> \
dbt("run_marts", "run --select tag:marts")
Обратите внимание: run_date прокидывается из контекста Airflow ({{ ds }}) в переменную dbt. Это и
есть склейка двух миров — оркестратор владеет временем, dbt владеет логикой.
7. Backfill: пересчёт прошлого
Backfill — заполнение исторических интервалов: либо потому что пайплайн только создан и нужна история за два года, либо потому что нашли баг и старые цифры неверны.
Backfill безопасен ровно тогда, когда выполнены три условия из картинки: непересекающиеся области записи, идемпотентность, ограниченный параллелизм. Если хоть одно нарушено — backfill не восстановит данные, а испортит их ещё сильнее.
Механика в Airflow
# Пересчитать март 2026 целиком (конец интервала исключается)
airflow dags backfill orders_daily --start-date 2026-03-01 --end-date 2026-03-31 --reset-dagruns
# Точечно: сбросить витрину и всё, что ниже по графу, за один день
airflow tasks clear orders_daily --task-regex "load_fct_orders" --downstream \
--start-date 2026-03-14 --end-date 2026-03-14
clear — самая частая операция дежурного: сброс состояния задач заставляет планировщик выполнить их
заново. Флаг --downstream почти всегда нужен: если изменился факт, витрины на его основе тоже неверны.
Стратегии и цена
- Хронологический (по одному дню). Обязателен, если у пайплайна есть состояние (накопительный итог, SCD2): день N+1 читает результат дня N.
- Параллельный. Дни независимы → пачками по 10–20 с ограничением пулом.
- Одним большим запросом. Часто самый быстрый путь: вместо 730 запусков — один SQL с
GROUP BY dtза весь период. Оркестратор тут не нужен; это разовая миграция или задачаfull_refreshс ручным триггером.
Оценку стоимости почти никогда не делают заранее, а она тривиальна:
время ≈ (кол-во_интервалов / параллелизм) × длительность_запуска. 730 дней по 12 минут при
параллелизме 1 — 6 суток; при параллелизме 20 — 7 часов. Разница определяется ровно одним
фактором: наличием состояния между днями. Поэтому «делайте дни независимыми» — не эстетика, а прямая
экономия суток дежурства.
Ловушка: измерение «на сегодня» вместо «на тот момент»
Пересчитывая март 2026 в июле, вы читаете dim_customer сегодняшний, а не мартовский. Сегменты
клиентов с тех пор изменились — и мартовская выручка по сегментам «поедет». Это не баг оркестратора,
а вопрос моделирования: измерение должно быть версионированным (SCD2), а джойн — по интервалу
валидности (f.order_ts >= d.valid_from and f.order_ts < d.valid_to). Подробности —
в https://courses.digitable.life/post/data-engineering/02-data-modeling/.
8. Как работает планировщик: последовательность событий
Чтобы отлаживать «почему задача не стартует», полезно держать в голове точную последовательность:
а исчерпанный пул end end Note over S,DB: если heartbeat пропал > threshold —
задача помечается zombie и уходит в retry
Отсюда прямо читаются три самых частых диагноза:
- Задача в
scheduled— нет слота (пул,max_active_tasks,parallelism), а не «Airflow завис». - Задача в
queued— нет живого воркера или очередь исполнителя не совпадает с очередью задачи. - DAG не появился в UI — DAG Processor не смог разобрать файл (ошибка импорта, таймаут парсинга).
9. Выбор инструмента: trade-offs
Коротко о позициях:
- cron + скрипты. Честный выбор для 5 задач. Ломается на первой зависимости между машинами: нет графа, ретраев, истории, backfill. «Переросли» наступает примерно на 20 джобах.
- Airflow. Зрелость и огромная экосистема провайдеров — любая интеграция уже написана. Цена — эксплуатация: метабаза, воркеры, задержка планировщика, тяжёлое тестирование DAG’ов.
- Dagster. Строит граф вокруг ассетов (таблиц), а не задач: линидж получается бесплатно. Ближе к тому, как мыслят дата-инженеры; экосистема меньше.
- Prefect. Минимум церемоний, обычный Python, динамические графы — когда структура пайплайна известна только в рантайме.
- Temporal. Не дата-оркестратор, а движок долгоживущих workflow с durable execution: саги и месяцами живущие процессы. Для витрин избыточен.
- dbt. Не альтернатива, а слой ниже: почти всегда живёт внутри одного из перечисленных.
Отдельно: Airflow 3 (2025) сдвинул модель в сторону ассетов и убрал прямой доступ воркеров к метабазе в пользу Task Execution API — разница между «задачным» и «ассетным» подходом стирается (release notes).
10. Типичные ошибки
| Ошибка | Что происходит | Лечение |
|---|---|---|
| Тяжёлый код на верхнем уровне DAG-файла | запрос к базе выполняется при каждом парсинге — сотни раз в час, парсинг падает по таймауту | структура из статики/Variable, никаких обращений вовне |
@daily при задержке источника |
данные приезжают в 02:40, DAG стартует в 00:00 и считает пустоту | сдвинуть расписание (0 3 * * *) или перейти на датасеты |
| Ретраи без идемпотентности | упали после вставки половины строк, ретрай вставил их снова | delete+insert окна или MERGE (раздел 3) |
| Ретраи без backoff | бьём по лежащему источнику каждые 10 с × 200 задач | retry_exponential_backoff=True, max_retry_delay |
| Данные через XCom | DataFrame в метабазе убивает Postgres оркестратора | XCom хранит путь в S3, данные едут мимо |
| Нет таймаутов | зависшая задача держит слот бесконечно, за ней встаёт очередь | execution_timeout на задаче, dagrun_timeout на DAG |
Ветвление по now() |
«если сегодня понедельник» — при backfill логика разъезжается | ветвиться по logical_date |
| Один DAG на всю компанию | 1200 задач, минуты на планирование, ревью невозможен | резать по владельцу и SLA, связывать датасетами |
Случайный catchup=True при старом start_date |
первый деплой порождает 730 запусков и кладёт источник | осознанный catchup, max_active_runs, пулы |
| Секреты в коде DAG | пароль утекает в git, логи и UI | Connections и секрет-бэкенды (Vault, Secrets Manager) |
11. Как это выглядит в проде
SLA и алерты по смыслу, а не по факту. Алерт «задача упала» бесполезен, если у неё 3 ретрая и она починится сама. Полезные сигналы: «витрина не обновилась к 09:00», «время выполнения выросло вдвое к скользящему среднему», «строк на 40% меньше обычного». Это метрики данных, а не процессов.
Проверка целостности DAG’ов в CI ловит половину инцидентов ещё до деплоя:
# tests/test_dag_integrity.py — обязательный тест в любом репозитории DAG'ов
import pytest
from airflow.models import DagBag
DAGBAG = DagBag(include_examples=False)
def test_no_import_errors():
"""DAG-файлы парсятся без ошибок и не содержат циклов."""
assert not DAGBAG.import_errors, f"Ошибки импорта: {DAGBAG.import_errors}"
@pytest.mark.parametrize("dag_id", DAGBAG.dag_ids)
def test_dag_conventions(dag_id):
"""Соглашения команды: владелец, теги, таймаут и ретраи на каждой задаче."""
dag = DAGBAG.get_dag(dag_id)
assert dag.default_args.get("owner"), f"{dag_id}: не указан owner"
assert dag.tags, f"{dag_id}: нет тегов"
for task in dag.tasks:
assert task.execution_timeout is not None, f"{dag_id}.{task.task_id}: нет таймаута"
assert task.retries >= 1, f"{dag_id}.{task.task_id}: нет ретраев"
Изоляция зависимостей. Разным пайплайнам нужны разные версии библиотек. Практика: каждая задача —
контейнер (KubernetesPodOperator), Airflow остаётся тонким слоем управления, и обновление pandas в
одном пайплайне не ломает остальные сорок.
Наблюдаемость платформы. Минимум дашбордов: длительность цикла планировщика, глубина очереди по
пулам, доля задач в queued дольше N минут, топ-20 задач по росту длительности, SLA-мисс за сутки.
Стоимость. Сам оркестратор дёшев (пара инстансов), но он генерирует расходы: каждый запуск — это compute в Spark/Snowflake. Пересчёт года «на всякий случай» может стоить дороже месячного бюджета команды, поэтому длинный backfill в зрелых командах проходит через явное согласование.
12. Мини-итог
- Оркестрация — про зависимости, время и восстановление после сбоев, а не про «запуск по крону».
- Абстракция — DAG; планирование сводится к топологической сортировке с ограничениями по ресурсам.
- Время в задачу инжектируется (
data_interval), никогда не берётся из системных часов. - Идемпотентность — фундамент. Без неё ретраи и backfill опасны, а не полезны. Три рабочих паттерна: overwrite партиции, MERGE по ключу, дедупликация по детерминированному ключу.
- Backfill возможен там, где области записи не пересекаются и параллелизм ограничен пулами.
- Airflow оркестрирует задачи, dbt — SQL-модели; связывает их проброс логической даты.
- Прод — это таймауты, пулы, backoff, тесты целостности DAG в CI и алерты по свежести данных.
Источники
- Apache Airflow — документация: DAGs, Scheduling, Pools, Deferrable Operators, а также Best Practices.
- dbt Developer Hub — инкрементальные модели и стратегии материализации.
- Dagster: Software-Defined Assets — ассетный подход к оркестрации.
- Martin Kleppmann. Designing Data-Intensive Applications, гл. 10–11 — батч, потоки, идемпотентность, происхождение данных. dataintensive.net
- Ralph Kimball, Margy Ross. The Data Warehouse Toolkit, 3rd ed. — SCD и правила загрузки измерений.
- Maxime Beauchemin. Functional Data Engineering — манифест идемпотентных, воспроизводимых пайплайнов.
- CLRS, Introduction to Algorithms, §22.4 — топологическая сортировка, теоретическая база планировщика.
Что дальше
Оркестратор умеет надёжно запускать пересчёты — но во что именно он пишет и почему одни форматы позволяют переписать одну партицию за секунды, а другие требуют перечитать всю таблицу? Разбираемся в следующей статье: Хранилища и форматы: Parquet, колоночные БД, lakehouse.