Data Engineering и ETL Оркестрация пайплайнов: Airflow, dbt, идемпотентность и backfill
0%

Оркестрация пайплайнов: Airflow, dbt, идемпотентность и backfill

Оркестрация пайплайнов: 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

Реальный планировщик усложняет это в трёх местах, и именно там живёт вся инженерия:

  1. Ресурсы. Готовых задач может быть 200, а слотов — 20. «Выбрать оптимальный порядок при ограниченных ресурсах» — это job-shop scheduling, NP-трудная задача, поэтому все оркестраторы используют жадные эвристики (приоритет, FIFO внутри приоритета) и пулы.
  2. Отказы. Задача не просто «выполнена/не выполнена» — есть ретраи, таймауты, зомби-процессы.
  3. Время. Граф запускается не один раз, а по расписанию: получается сетка «граф × момент времени». Именно из-за этого измерения появляются backfill и идемпотентность.

Вот тот же пайплайн как граф — обратите внимание на естественный параллелизм двух веток:

Важная деталь: 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. Жизненный цикл задачи

Задача — не «выполнена / не выполнена». В любом зрелом оркестраторе это конечный автомат, и понимание его состояний экономит часы отладки:

Три состояния тут неочевидны:

  • queuedrunning. Задача может часами висеть в 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 безопасен ровно тогда, когда выполнены три условия из картинки: непересекающиеся области записи, идемпотентность, ограниченный параллелизм. Если хоть одно нарушено — 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. Как работает планировщик: последовательность событий

Чтобы отлаживать «почему задача не стартует», полезно держать в голове точную последовательность:

Отсюда прямо читаются три самых частых диагноза:

  1. Задача в scheduledнет слота (пул, max_active_tasks, parallelism), а не «Airflow завис».
  2. Задача в queuedнет живого воркера или очередь исполнителя не совпадает с очередью задачи.
  3. 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.

Нашли неточность? Выделите фрагмент текста — рядом появится жучок.

Нужен разбор именно вашей ситуации?

Статья описывает общий случай. Если у вас частный — можно разобрать его отдельно, платно. А если не хватает целого материала, предложите тему: её оплачивают вскладчину, и она выходит открытой для всех.

Доска запросов