dbt + Airflow: BashOperator dbt run vs SparkSubmitOperator, CI/CD pipeline

dbt + Airflow: BashOperator dbt run vs SparkSubmitOperator, CI/CD pipeline

platform

Архитектура оркестрации Modern Data Lakehouse

Почему оркестрация - отдельный слой

Когда дата-инженер начинает работать с dbt, возникает соблазн думать, что dbt делает всё: он управляет зависимостями, строит граф, запускает трансформации. Зачем тогда Airflow?

dbt действительно управляет зависимостями - но только внутри своего проекта. dbt не умеет:

  • Запускаться по расписанию (cron, event-driven)
  • Ждать, пока завершится загрузка данных из внешнего источника
  • Запускать не-dbt задачи: Python-скрипты, Spark-джобы, ML-инференс
  • Отслеживать SLA: «если витрина не обновилась к 8:00 - отправить алерт»
  • Повторно запускать упавшую задачу через 5 минут
  • Строить зависимости между разными системами (dbt → ML-pipeline → отчёт)

Airflow закрывает всё это. В modern data stack принято разделение ответственности:

  • Airflow - дирижёр: расписание, координация, мониторинг, алертинг
  • dbt - трансформации: SQL-модели, тесты, документация
  • Spark - вычислительный движок: физическое выполнение SQL

Важно понимать топологию: когда Airflow запускает dbt run, физические вычисления не происходят на Airflow-воркере. Airflow-воркер лишь отправляет команду в dbt CLI, dbt компилирует SQL и через адаптер отправляет его на Spark Thrift Server или Databricks API. Вся нагрузка на CPU и память ложится на кластер Spark - Airflow-воркер в это время просто ждёт ответа.

Airflow Worker           dbt CLI              Spark Cluster
    │                      │                      │
    ├─ dbt run ──────────► │                      │
    │                      ├─ compile SQL ────►   │
    │                      ├─ send to Thrift ──►  │
    │                      │                      ├─ execute
    │                      │                      ├─ return result
    │                      │◄─ OK / ERROR ─────── │
    │◄─ task success ───── │                      │

Airflow-воркер потребляет минимум ресурсов - он занимается координацией, а не вычислениями. Это важно для правильного планирования инфраструктуры.

Проблема «толстого» DAG-файла

Один из типичных anti-patterns - попытка превратить каждую dbt-модель в отдельный Airflow-таск через динамическую генерацию DAG. Выглядит логично: тогда Airflow видит каждую модель отдельно, можно перезапустить конкретную упавшую модель.

Проблема в том, что Airflow Scheduler непрерывно парсит все DAG-файлы. Если DAG-файл читает manifest.json (артефакт dbt с описанием всех моделей) и создаёт тысячи Python-объектов Task - это происходит при каждом парсинге, каждые несколько секунд. При проекте из 500+ моделей это может перегрузить Scheduler.

Решения:

  1. Не парсить manifest.json в DAG-файле - задать модели вручную или группами (по тегам)
  2. Использовать Cosmos - специализированная библиотека, которая делает это эффективно
  3. Использовать group-based подход - один таск на слой (staging, marts)

Роль Airflow Scheduler и Executor

Понимание компонентов Airflow помогает правильно отлаживать проблемы:

  • Scheduler - непрерывно парсит DAG-файлы, находит задачи к запуску, ставит в очередь
  • Executor - определяет, где физически выполняется таск: LocalExecutor (тот же процесс), CeleryExecutor (воркеры-воркеры), KubernetesExecutor (Pod в K8s)
  • Worker - физический процесс, который выполняет Python-код таска (например, BashOperator.execute())
  • Webserver - UI для мониторинга, не участвует в выполнении
  • Metadata DB - PostgreSQL/MySQL, хранит состояние тасков, DAG-раны, логи

При использовании BashOperator с dbt run: Worker запускает shell-процесс dbt run, ждёт его завершения, записывает stdout/stderr в логи Airflow. Весь этот процесс происходит на машине Worker - то есть на ней должен быть установлен dbt и настроен profiles.yml.


BashOperator + dbt run: простота и практичность

Механика BashOperator

BashOperator - это самый прямолинейный способ запустить dbt из Airflow. Airflow-воркер выполняет shell-команду, захватывает stdout/stderr и записывает в логи. Возвращаемый код (exit code) определяет успех или провал таска: 0 = OK, любой другой = FAILED.

dbt CLI возвращает ненулевой exit code при любой ошибке: упавший тест, синтаксическая ошибка в SQL, недостижимый Spark-кластер. Это делает интеграцию естественной - Airflow автоматически пометит таск как упавший.

Простой DAG с BashOperator

# dags/dbt_transactions_pipeline.py
from datetime import datetime, timedelta
from airflow import DAG
from airflow.operators.bash import BashOperator

# Константы конфигурации
DBT_PROJECT_DIR = "/opt/dbt/my_project"
DBT_PROFILES_DIR = "/opt/dbt"
DBT_TARGET = "prod"

# Переменные окружения для dbt - передаются в shell-процесс
DBT_ENV = {
    "DBT_TARGET": DBT_TARGET,
    "SPARK_HOST": "{{ var.value.spark_thrift_host }}",
    "SPARK_USER": "airflow_svc",
    "SPARK_PASSWORD": "{{ conn.spark_default.password }}",
    "DBT_LOCATION_ROOT": "s3a://prod-lake/gold",
}

default_args = {
    "owner": "data-platform-team",
    "depends_on_past": False,
    "email": ["data-alerts@company.ru"],
    "email_on_failure": True,
    "email_on_retry": False,
    "retries": 2,
    "retry_delay": timedelta(minutes=5),
    "retry_exponential_backoff": True,  # 5min → 10min → 20min
}

with DAG(
    dag_id="dbt_transactions_pipeline",
    description="Ежедневный пересчёт финансовых витрин через dbt + Spark",
    schedule="0 6 * * *",          # Каждый день в 06:00 UTC
    start_date=datetime(2024, 1, 1),
    catchup=False,                  # Не запускать пропущенные раны
    max_active_runs=1,             # Только один активный ран одновременно
    tags=["dbt", "spark", "finance"],
    default_args=default_args,
) as dag:

    # Шаг 1: Проверить свежесть источников
    check_freshness = BashOperator(
        task_id="check_source_freshness",
        bash_command=(
            f"cd {DBT_PROJECT_DIR} && "
            f"dbt source freshness "
            f"--profiles-dir {DBT_PROFILES_DIR} "
            f"--target {DBT_TARGET} "
            f"--select source:bronze_transactions"
        ),
        env=DBT_ENV,
    )

    # Шаг 2: Запустить staging-слой с тестами
    build_staging = BashOperator(
        task_id="build_staging",
        bash_command=(
            f"cd {DBT_PROJECT_DIR} && "
            f"dbt build "
            f"--profiles-dir {DBT_PROFILES_DIR} "
            f"--target {DBT_TARGET} "
            f"--select tag:staging"
        ),
        env=DBT_ENV,
    )

    # Шаг 3: Запустить marts-слой с тестами
    build_marts = BashOperator(
        task_id="build_marts",
        bash_command=(
            f"cd {DBT_PROJECT_DIR} && "
            f"dbt build "
            f"--profiles-dir {DBT_PROFILES_DIR} "
            f"--target {DBT_TARGET} "
            f"--select tag:finance_marts"
        ),
        env=DBT_ENV,
    )

    # Шаг 4: Генерация документации (опционально, в prod)
    generate_docs = BashOperator(
        task_id="generate_docs",
        bash_command=(
            f"cd {DBT_PROJECT_DIR} && "
            f"dbt docs generate "
            f"--profiles-dir {DBT_PROFILES_DIR} "
            f"--target {DBT_TARGET}"
        ),
        env=DBT_ENV,
        # Падение docs не должно ломать весь DAG
        trigger_rule="all_success",
    )

    # Определяем порядок выполнения
    check_freshness >> build_staging >> build_marts >> generate_docs

Этот DAG делает четыре вещи последовательно: проверяет свежесть источников, строит staging с тестами, строит marts с тестами, генерирует документацию. Если какой-то шаг упал - следующие не запускаются.

dbt selectors: точечный запуск моделей

Флаг --select в dbt - мощный механизм выбора подмножества моделей. Понимание селекторов критично для эффективного DAG-дизайна:

# По тегу
dbt build --select tag:finance_marts

# По пути (директория)
dbt build --select models/marts/finance/

# Конкретная модель
dbt build --select fct_transactions

# Модель и все её downstream зависимости
dbt build --select fct_transactions+

# Модель и все её upstream зависимости
dbt build --select +fct_transactions

# Изменённые модели (относительно prod manifest)
dbt build --select state:modified+

# Комбинация: изменённые модели тега finance
dbt build --select "tag:finance AND state:modified+"

# Исключение
dbt build --select tag:finance --exclude tag:slow_tests

В Airflow-DAG используйте теги для группировки моделей по слоям и доменам. Это позволяет перезапустить только нужный слой, не трогая остальные.

Передача credentials через Airflow

Хранить пароли в DAG-файле нельзя. Airflow предоставляет несколько механизмов:

Airflow Variables - простые key-value хранилище:

from airflow.models import Variable

spark_host = Variable.get("spark_thrift_host")

Airflow Connections - структурированные credentials (host, port, login, password, extra):

from airflow.hooks.base import BaseHook

conn = BaseHook.get_connection("spark_thrift_prod")
spark_host = conn.host
spark_password = conn.password

Secrets Backend (production-рекомендация) - Airflow читает secrets из HashiCorp Vault или AWS Secrets Manager:

# airflow.cfg или docker-compose переменная
[secrets]
backend = airflow.providers.hashicorp.secrets.vault.VaultBackend
backend_kwargs = {"connections_path": "airflow/connections", "variables_path": "airflow/variables"}

В DAG передаём credentials как переменные окружения в shell-процесс dbt:

from airflow.hooks.base import BaseHook

def get_dbt_env():
    conn = BaseHook.get_connection("spark_thrift_prod")
    return {
        "SPARK_HOST": conn.host,
        "SPARK_PORT": str(conn.port),
        "SPARK_USER": conn.login,
        "SPARK_PASSWORD": conn.password,
        "DBT_LOCATION_ROOT": Variable.get("dbt_location_root_prod"),
    }

build_marts = BashOperator(
    task_id="build_marts",
    bash_command="cd /opt/dbt && dbt build --select tag:finance_marts",
    env=get_dbt_env(),
)

В profiles.yml dbt читает эти переменные через env_var():

# profiles.yml
my_project:
  target: "{{ env_var('DBT_TARGET', 'dev') }}"
  outputs:
    prod:
      type: spark
      method: thrift
      host: "{{ env_var('SPARK_HOST') }}"
      port: "{{ env_var('SPARK_PORT', '10000') | int }}"
      user: "{{ env_var('SPARK_USER') }}"
      password: "{{ env_var('SPARK_PASSWORD') }}"
      schema: gold
      threads: 8

dbt build vs dbt run в Airflow-DAG

Критически важный выбор при построении DAG:

# ПЛОХО: run + test раздельно - тесты не блокируют downstream
run_task = BashOperator(bash_command="dbt run --select +fct_orders")
test_task = BashOperator(bash_command="dbt test --select +fct_orders")
run_task >> test_task >> downstream_mart  # fct_orders уже в Gold, даже если тест упадёт

# ХОРОШО: build = run + test в правильном порядке по DAG
build_task = BashOperator(bash_command="dbt build --select +fct_orders")
build_task >> downstream_mart  # downstream не запустится, если test упал

dbt build гарантирует, что downstream-модели получают только валидированные данные. Всегда используйте dbt build в production-DAG.

Граф зависимостей: несколько доменов

Для крупных проектов DAG разбивают по доменам:

# dags/dbt_full_pipeline.py
with DAG(dag_id="dbt_full_pipeline", schedule="0 5 * * *", ...) as dag:

    freshness = BashOperator(
        task_id="source_freshness",
        bash_command="cd /opt/dbt && dbt source freshness",
    )

    staging = BashOperator(
        task_id="build_staging",
        bash_command="cd /opt/dbt && dbt build --select tag:staging",
    )

    intermediate = BashOperator(
        task_id="build_intermediate",
        bash_command="cd /opt/dbt && dbt build --select tag:intermediate",
    )

    # Параллельное выполнение финансового и маркетингового домена
    finance = BashOperator(
        task_id="build_finance_marts",
        bash_command="cd /opt/dbt && dbt build --select tag:finance",
    )

    marketing = BashOperator(
        task_id="build_marketing_marts",
        bash_command="cd /opt/dbt && dbt build --select tag:marketing",
    )

    ml_features = BashOperator(
        task_id="build_ml_features",
        bash_command="cd /opt/dbt && dbt build --select tag:ml_features",
    )

    # Граф зависимостей
    freshness >> staging >> intermediate >> [finance, marketing]
    finance >> ml_features

Airflow выполнит finance и marketing параллельно (при наличии свободных воркеров), но оба стартуют только после успешного intermediate.

Ограничения BashOperator

BashOperator прост, но имеет ограничения:

  • Нет fine-grained control: если dbt-модель внутри dbt build --select tag:staging упала, Airflow видит только «таск build_staging упал». Нельзя перезапустить конкретную модель.
  • Нет нативной интеграции с Spark UI: Airflow-логи показывают вывод dbt CLI, но не Spark Application ID. Чтобы найти джоб в Spark UI - нужно искать вручную.
  • dbt должен быть установлен на Worker: каждый Airflow-воркер должен иметь dbt + dbt-spark + правильный profiles.yml. При обновлении dbt нужно обновить все воркеры.

SparkSubmitOperator: нативный Spark без dbt

Когда нужен SparkSubmitOperator

SparkSubmitOperator - это оператор из apache-airflow-providers-apache-spark, который запускает Spark-приложения через spark-submit. Он не связан с dbt и предназначен для другого класса задач:

  • Сложные ETL-операции, которые нельзя выразить в SQL (машинное обучение, граф-алгоритмы, кастомные парсеры)
  • Интеграция с данными, которые dbt не поддерживает (бинарные форматы, RDD-операции)
  • Тяжёлые batch-джобы с тонкой настройкой ресурсов (specific executor memory, custom JVM flags)
from airflow.providers.apache.spark.operators.spark_submit import SparkSubmitOperator

run_spark_etl = SparkSubmitOperator(
    task_id="run_spark_etl",
    application="/opt/spark-apps/etl_raw_events.py",
    conn_id="spark_default",             # Airflow Connection с URL кластера
    name="airflow_etl_raw_events",
    deploy_mode="cluster",               # cluster: driver на воркере кластера
    num_executors=10,
    executor_memory="4g",
    executor_cores=2,
    driver_memory="2g",
    jars="/opt/spark-jars/delta-core.jar,/opt/spark-jars/iceberg-spark.jar",
    conf={
        "spark.sql.shuffle.partitions": "200",
        "spark.sql.adaptive.enabled": "true",
        "spark.sql.extensions": "io.delta.sql.DeltaSparkSessionExtension",
    },
    application_args=[
        "--date", "{{ ds }}",            # Airflow templating: дата запуска
        "--output", "s3a://bronze/events",
    ],
)

deploy_mode: client vs cluster

Это важный параметр, влияющий на топологию выполнения:

Client mode (deploy_mode="client"): Spark Driver запускается на машине, которая вызвала spark-submit - то есть на Airflow-воркере. Executors запускаются на кластере. Если воркер упадёт - Driver умрёт, джоб упадёт.

Cluster mode (deploy_mode="cluster"): Spark Driver запускается на случайном узле кластера. Airflow-воркер только инициирует запуск и может отключиться. Driver живёт независимо от Airflow. Правильный выбор для production.

# Client mode - для разработки и дебага:
run_etl = SparkSubmitOperator(
    ...,
    deploy_mode="client",  # Driver на воркере, легче читать логи
)

# Cluster mode - для production:
run_etl = SparkSubmitOperator(
    ...,
    deploy_mode="cluster",  # Driver на кластере, независим от Airflow
)

Гибридный паттерн: Spark ETL + dbt

Реальный production-пайплайн часто комбинирует оба подхода:

with DAG("hybrid_spark_dbt_pipeline", schedule="0 4 * * *", ...) as dag:

    # Шаг 1: Сырой Spark-джоб для парсинга нестандартных форматов
    parse_raw_logs = SparkSubmitOperator(
        task_id="parse_raw_logs",
        application="/apps/parse_nginx_logs.py",
        conn_id="spark_prod",
        deploy_mode="cluster",
        executor_memory="8g",
        num_executors=20,
        application_args=["--date", "{{ ds }}"],
    )

    # Шаг 2: Spark-джоб для ML-фичей (нельзя в dbt)
    compute_embeddings = SparkSubmitOperator(
        task_id="compute_user_embeddings",
        application="/apps/user_embeddings.py",
        conn_id="spark_prod",
        deploy_mode="cluster",
        executor_memory="16g",
        num_executors=30,
        conf={"spark.ml.dlc.enabled": "true"},
        application_args=["--date", "{{ ds }}"],
    )

    # Шаг 3: dbt строит витрины поверх результатов Spark
    build_marts = BashOperator(
        task_id="build_marketing_marts",
        bash_command="cd /opt/dbt && dbt build --select tag:marketing",
        env=get_dbt_env(),
    )

    # Шаг 4: dbt строит ML-фичи поверх embeddings
    build_ml_features = BashOperator(
        task_id="build_ml_features",
        bash_command="cd /opt/dbt && dbt build --select tag:ml",
        env=get_dbt_env(),
    )

    # Граф: Spark-джобы параллельно, затем dbt поверх
    [parse_raw_logs, compute_embeddings] >> build_marts
    compute_embeddings >> build_ml_features

Здесь Spark делает то, что не умеет dbt (парсинг нестандартных форматов, ML), а dbt делает то, что не стоит писать на PySpark (SQL-трансформации, тесты, документация).


Cosmos: автоматический парсинг dbt DAG

Что такое Cosmos и зачем он нужен

Astronomer Cosmos (apache-airflow-providers-astronomer-cosmos) - это библиотека, которая решает главную проблему BashOperator: непрозрачность. С Cosmos каждая dbt-модель становится отдельным Airflow-таском. Провал одной модели не убивает весь DAG - можно перезапустить только её.

Cosmos читает manifest.json (скомпилированный артефакт dbt) и строит Airflow TaskGroup, точно воспроизводя граф зависимостей dbt.

Установка и базовая конфигурация

pip install astronomer-cosmos[dbt-spark]
# dags/dbt_cosmos_pipeline.py
from cosmos import DbtTaskGroup, ProjectConfig, ProfileConfig, ExecutionConfig, RenderConfig
from cosmos.profiles import SparkThriftProfileMapping
from airflow.decorators import dag
from datetime import datetime

@dag(
    schedule="0 6 * * *",
    start_date=datetime(2024, 1, 1),
    catchup=False,
    tags=["cosmos", "dbt", "spark"],
)
def dbt_cosmos_pipeline():

    # Описываем профиль подключения к Spark
    profile_config = ProfileConfig(
        profile_name="my_project",
        target_name="prod",
        profile_mapping=SparkThriftProfileMapping(
            conn_id="spark_thrift_prod",   # Airflow Connection
            profile_args={
                "schema": "gold",
                "threads": 8,
            },
        ),
    )

    # Где находится dbt-проект
    project_config = ProjectConfig(
        dbt_project_path="/opt/dbt/my_project",
    )

    # Настройки рендеринга TaskGroup
    render_config = RenderConfig(
        select=["tag:finance"],  # Только финансовые модели
        exclude=["tag:slow"],    # Кроме медленных тестов
    )

    # Создаём TaskGroup из всех dbt-моделей
    finance_group = DbtTaskGroup(
        group_id="finance_models",
        project_config=project_config,
        profile_config=profile_config,
        render_config=render_config,
    )

    finance_group

dbt_cosmos_pipeline()

Что видит инженер в Airflow UI

С Cosmos каждая модель - отдельный таск. Если fct_transactions упала, а fct_orders прошла успешно - можно перезапустить только fct_transactions. Все тесты модели тоже видны как отдельные таски.

finance_models/
  ├── stg_transactions            [success]
  ├── stg_orders                  [success]
  ├── int_transactions_enriched   [success]
  ├── fct_transactions            [failed]    ← перезапускаем только это
  ├── fct_transactions.not_null.transaction_id  [skipped]
  └── fct_orders                  [success]

Это принципиальный выигрыш по сравнению с BashOperator, где видно только «build_finance_marts: FAILED».

Когда Cosmos, когда BashOperator

Критерий BashOperator Cosmos
Размер проекта До 50 моделей 50+ моделей
Отладка Сложно: нет гранулярности Легко: каждая модель видна
Перезапуск Весь слой целиком Конкретная модель
Сложность настройки Минимальная Средняя
Нагрузка на Scheduler Минимальная Средняя (много тасков)
Зависимости Python dbt CLI dbt + Cosmos

Для начала рекомендуется BashOperator - он проще и понятнее. Cosmos стоит внедрять, когда проект вырос до нескольких сотен моделей и перезапуск всего слоя становится болезненным.


Архитектурная диаграмма: dbt + Airflow + Spark

Схема показывает полный путь от кода до данных: изменения в Git проходят через CI/CD, попадают на Airflow Scheduler, который запускает Worker. Worker вызывает dbt CLI, dbt через Thrift Server отправляет SQL на кластер. Spark выполняет вычисления и записывает результат в S3/MinIO, обновляя метаданные в Hive Metastore.


CI/CD пайплайн для аналитики

Концепция «Данные как код»

В современной аналитической инженерии принято относиться к SQL-моделям так же, как к коду приложений: никаких изменений без code review, никаких изменений без прохождения тестов, полная история в Git.

Это называется GitOps для данных (или «Data as Code»). Последствия:

  • Каждое изменение витрины - Pull Request с review от коллеги
  • CI автоматически проверяет корректность до merge
  • Деплой в production - автоматически после merge в main
  • Откат к предыдущей версии - git revert

Анатомия CI-пайплайна при Pull Request

# .github/workflows/ci.yml
name: dbt CI

on:
  pull_request:
    branches: [main]
    paths:
      - 'models/**'
      - 'tests/**'
      - 'macros/**'
      - 'dbt_project.yml'

env:
  DBT_TARGET: ci
  SPARK_HOST: ${{ secrets.SPARK_CI_HOST }}
  SPARK_USER: ci_svc
  SPARK_PASSWORD: ${{ secrets.SPARK_CI_PASSWORD }}
  DBT_LOCATION_ROOT: s3a://ci-lake/gold

jobs:
  dbt-ci:
    name: dbt CI checks
    runs-on: ubuntu-latest

    steps:
      - name: Checkout code
        uses: actions/checkout@v4

      # ──────────────────────────────────────────
      # Шаг 1: Линтинг SQL (синтаксис + стиль)
      # ──────────────────────────────────────────
      - name: Install SQLFluff
        run: pip install sqlfluff sqlfluff-templater-dbt==2.3.0

      - name: Lint SQL with SQLFluff
        run: |
          sqlfluff lint models/ \
            --dialect sparksql \
            --templater dbt \
            --config .sqlfluff \
            --exclude-rules L031,L034

      # ──────────────────────────────────────────
      # Шаг 2: Установка dbt
      # ──────────────────────────────────────────
      - name: Install dbt-spark
        run: pip install dbt-spark[PyHive]==1.7.4

      - name: Install dbt packages
        run: dbt deps

      # ──────────────────────────────────────────
      # Шаг 3: Compile - проверка Jinja и SQL без выполнения
      # ──────────────────────────────────────────
      - name: dbt compile
        run: dbt compile --target ci
        # Если компиляция упала - есть синтаксическая ошибка в Jinja или SQL

      # ──────────────────────────────────────────
      # Шаг 4: Slim CI - запуск только изменённых моделей
      # ──────────────────────────────────────────
      - name: Download prod manifest (для state:modified)
        run: |
          aws s3 cp s3://dbt-artifacts/prod/manifest.json \
            ./prod_manifest/manifest.json
        env:
          AWS_ACCESS_KEY_ID: ${{ secrets.AWS_ACCESS_KEY_ID }}
          AWS_SECRET_ACCESS_KEY: ${{ secrets.AWS_SECRET_ACCESS_KEY }}

      - name: dbt build (только изменённые модели)
        run: |
          dbt build \
            --target ci \
            --select state:modified+ \
            --defer \
            --state ./prod_manifest/
        # state:modified+ - изменённые модели + их downstream
        # --defer - для upstream моделей использовать prod данные
        # --state - где искать manifest.json предыдущего рана

Slim CI: магия --defer и state:modified

Это самая умная часть CI-пайплайна. Без --defer каждый CI-ран пересчитывал бы весь проект - терабайты данных за каждый PR. С --defer:

  • state:modified+ - запустить только модели, которые изменились в PR (плюс все их downstream зависимости)
  • --defer - для upstream моделей (которые PR не трогал) использовать данные из production, а не пересчитывать

Пример: разработчик изменил fct_transactions. --select state:modified+ запустит только fct_transactions и всё, что от неё зависит. Upstream модели (stg_transactions, int_transactions_enriched) будут взяты из production схемы через --defer. CI обработает 100 GB вместо 10 TB - за 10 минут вместо 2 часов.

# Что происходит с --defer:
# 1. dbt смотрит в prod manifest: stg_transactions существует в prod
# 2. Для stg_transactions использует prod.stg_transactions (не пересчитывает)
# 3. Для fct_transactions создаёт ci_schema.fct_transactions (пересчитывает)
# 4. Тесты запускаются на ci_schema.fct_transactions

Анатомия CD-пайплайна при merge в main

# .github/workflows/cd.yml
name: dbt CD

on:
  push:
    branches: [main]

jobs:
  dbt-deploy:
    name: Deploy to production
    runs-on: ubuntu-latest

    steps:
      - uses: actions/checkout@v4

      - name: Install dbt
        run: pip install dbt-spark[PyHive]==1.7.4

      - name: Install dbt packages
        run: dbt deps

      # ──────────────────────────────────────────
      # Компилируем проект (создаёт manifest.json)
      # ──────────────────────────────────────────
      - name: dbt compile (generate manifest)
        run: dbt compile --target prod
        env:
          DBT_TARGET: prod
          SPARK_HOST: ${{ secrets.SPARK_PROD_HOST }}
          # ... остальные секреты

      # ──────────────────────────────────────────
      # Загружаем manifest в S3 (для Slim CI следующих PR)
      # ──────────────────────────────────────────
      - name: Upload manifest to S3
        run: |
          aws s3 cp target/manifest.json \
            s3://dbt-artifacts/prod/manifest.json
        env:
          AWS_ACCESS_KEY_ID: ${{ secrets.AWS_ACCESS_KEY_ID }}
          AWS_SECRET_ACCESS_KEY: ${{ secrets.AWS_SECRET_ACCESS_KEY }}

      # ──────────────────────────────────────────
      # Синхронизируем dbt-проект на Airflow-воркеры
      # ──────────────────────────────────────────
      - name: Sync dbt project to Airflow workers
        run: |
          aws s3 sync . s3://airflow-dbt-projects/my_project/ \
            --exclude ".git/*" \
            --exclude "target/*" \
            --exclude "logs/*"
        env:
          AWS_ACCESS_KEY_ID: ${{ secrets.AWS_ACCESS_KEY_ID }}
          AWS_SECRET_ACCESS_KEY: ${{ secrets.AWS_SECRET_ACCESS_KEY }}

      # ──────────────────────────────────────────
      # Опционально: триггернуть Airflow DAG
      # ──────────────────────────────────────────
      - name: Trigger Airflow DAG
        run: |
          curl -X POST \
            -H "Content-Type: application/json" \
            -H "Authorization: Basic ${{ secrets.AIRFLOW_API_TOKEN }}" \
            -d '{"conf": {"triggered_by": "cd_pipeline", "commit": "${{ github.sha }}"}}' \
            ${{ secrets.AIRFLOW_URL }}/api/v1/dags/dbt_transactions_pipeline/dagRuns

После merge в main: компилируем проект, сохраняем manifest.json для Slim CI будущих PR, синхронизируем dbt-проект на Airflow-воркеры, опционально триггерим немедленный ран (или ждём следующего расписания).

Доставка dbt-проекта на Airflow-воркеры

Существует несколько стратегий деплоя dbt-проекта на воркеры:

S3-синхронизация (как выше): воркер при запуске таска скачивает проект из S3. Простой, но немного медленный.

GitSync sidecar (Kubernetes): Init-контейнер перед каждым запуском Pod'а делает git pull. Проект всегда актуален.

Docker-образ: dbt-проект запекается в Docker-образ. При обновлении CI собирает новый образ и пушит в registry. Airflow использует DockerOperator или KubernetesPodOperator с новым тегом образа.

# Вариант с DockerOperator - dbt запускается в изолированном контейнере
from airflow.providers.docker.operators.docker import DockerOperator

build_marts = DockerOperator(
    task_id="build_marts",
    image="registry.company.ru/dbt-project:{{ var.value.dbt_image_tag }}",
    command="dbt build --target prod --select tag:finance",
    environment={
        "SPARK_HOST": "{{ var.value.spark_host }}",
        "SPARK_PASSWORD": "{{ conn.spark_prod.password }}",
    },
    docker_url="unix://var/run/docker.sock",
    network_mode="host",
    auto_remove=True,
)

Docker-образ - лучшая стратегия для production: изоляция зависимостей, воспроизводимость, простой rollback (поменять тег образа).


Изоляция сред: dev / ci / prod

Три среды и их характеристики

Параметр dev ci prod
Spark cluster Локальный / shared dev Отдельный CI Production
Schema prefix dev_{username}_ ci_{pr_number}_ (без префикса)
Data volume Маленькая выборка state:modified+ Полный объём
Location s3a://dev-lake/... s3a://ci-lake/... s3a://prod-lake/...
Запуск Вручную / dbt CLI GitHub Actions Airflow
Freshness checks Пропускаем Запускаем Обязательно

profiles.yml для трёх сред

# profiles.yml
my_project:
  target: "{{ env_var('DBT_TARGET', 'dev') }}"

  outputs:

    dev:
      type: spark
      method: session   # Локальный SparkSession
      schema: "dev_{{ env_var('USER', 'unknown') }}"
      config:
        spark.master: "local[4]"
        spark.sql.adaptive.enabled: "true"

    ci:
      type: spark
      method: thrift
      host: "{{ env_var('SPARK_HOST') }}"
      port: 10000
      user: ci_svc
      password: "{{ env_var('SPARK_PASSWORD') }}"
      schema: "ci_{{ env_var('PR_NUMBER', 'local') }}"
      threads: 4

    prod:
      type: spark
      method: thrift
      host: "{{ env_var('SPARK_HOST') }}"
      port: 10000
      user: "{{ env_var('SPARK_USER', 'airflow_svc') }}"
      password: "{{ env_var('SPARK_PASSWORD') }}"
      schema: gold
      threads: 16
      connect_timeout: 120
      connect_retries: 3
      server_side_parameters:
        "spark.sql.shuffle.partitions": "400"
        "spark.sql.adaptive.enabled": "true"

Branch-based схемы для изоляции разработчиков

В dev-среде каждый разработчик работает в своей схеме - нет конфликтов при параллельной разработке:

# Разработчик Ivan
DBT_TARGET=dev USER=ivan dbt run --select fct_transactions
# Создаст: dev_ivan.fct_transactions

# Разработчик Maria работает параллельно
DBT_TARGET=dev USER=maria dbt run --select fct_transactions
# Создаст: dev_maria.fct_transactions

В CI схема изолирована по номеру PR:

# PR #42
DBT_TARGET=ci PR_NUMBER=42 dbt build --select state:modified+
# Создаст: ci_42.fct_transactions

После закрытия PR - схема ci_42 удаляется (можно автоматизировать через post-close workflow в GitHub Actions).


Бэкфилл: исторический пересчёт в Airflow

Проблема бэкфилла

Иногда нужно пересчитать исторические данные: исправлена бизнес-логика, изменилась схема источника, обнаружена ошибка в модели. При партиционированных таблицах бэкфилл означает перебор всех партиций.

dbt backfill через Airflow

# dags/dbt_backfill.py
from airflow.decorators import dag, task
from airflow.operators.bash import BashOperator
from datetime import datetime, timedelta

@dag(
    dag_id="dbt_backfill_fct_transactions",
    schedule=None,          # Только ручной запуск
    start_date=datetime(2024, 1, 1),
    catchup=False,
    tags=["backfill", "manual"],
    params={
        "start_date": "2023-01-01",
        "end_date": "2024-01-01",
        "model": "fct_transactions",
    }
)
def backfill_dag():

    backfill = BashOperator(
        task_id="dbt_backfill",
        bash_command="""
            cd /opt/dbt && dbt run \
                --select {{ params.model }} \
                --full-refresh \
                --vars '{"backfill_start": "{{ params.start_date }}", "backfill_end": "{{ params.end_date }}"}'
        """,
        env=get_dbt_env(),
    )

    backfill

backfill_dag()

В модели используем переменные для фильтрации:

-- models/marts/fct_transactions.sql
{% set backfill_start = var('backfill_start', none) %}
{% set backfill_end = var('backfill_end', none) %}

{{ config(
    materialized='incremental',
    file_format='delta',
    incremental_strategy='insert_overwrite'
) }}

SELECT * FROM {{ ref('int_transactions_enriched') }}

{% if is_incremental() and not backfill_start %}
    WHERE event_date > (SELECT MAX(event_date) FROM {{ this }})
{% elif backfill_start %}
    WHERE event_date BETWEEN '{{ backfill_start }}' AND '{{ backfill_end }}'
{% endif %}

--full-refresh + переменные диапазона дат - безопасный способ бэкфилла: модель пересоздаётся для заданного диапазона, исторические данные за пределами диапазона не трогаются.


Мониторинг и наблюдаемость

Spark Application ID в логах Airflow

Одна из главных проблем при отладке: Airflow показывает только вывод dbt CLI, а не Spark Application ID. Чтобы найти джоб в Spark UI - нужно искать вручную по времени.

Решение - передавать Application ID в логи явно. dbt поддерживает хуки через on-run-end:

# dbt_project.yml
on-run-end:
  - "{{ log_spark_application_id() }}"
{# macros/logging.sql #}
{% macro log_spark_application_id() %}
    {%- set result = run_query("SELECT spark_application_id()") -%}
    {%- if result -%}
        {{ log("SPARK_APP_ID: " ~ result.columns[0].values()[0], info=True) }}
    {%- endif -%}
{% endmacro %}

В логах Airflow появится строка SPARK_APP_ID: application_1234567890_0042 - теперь можно сразу перейти в Spark History Server.

SLA Miss мониторинг

Airflow поддерживает SLA-алерты: если таск не завершился к определённому времени - отправить уведомление:

from datetime import timedelta
from airflow.models import DAG

def on_sla_miss(dag, task_list, blocking_task_list, slas, blocking_tis):
    """Вызывается при пропуске SLA."""
    import requests
    message = f"SLA MISS: DAG {dag.dag_id}, tasks: {task_list}"
    # Отправить в Slack, PagerDuty, Telegram...
    requests.post(
        SLACK_WEBHOOK_URL,
        json={"text": f"⚠️ {message}"}
    )

with DAG(
    dag_id="dbt_transactions_pipeline",
    sla_miss_callback=on_sla_miss,
    default_args={
        "sla": timedelta(hours=2),  # SLA: таск должен завершиться за 2 часа
    },
) as dag:
    ...

dbt артефакты как источник метрик

dbt создаёт несколько JSON-артефактов после каждого рана:

  • target/run_results.json - результаты тасков: статус, время выполнения, количество строк
  • target/manifest.json - описание всего проекта: модели, тесты, зависимости
  • target/sources.json - результаты freshness-проверок

Эти файлы можно собирать и визуализировать в Grafana:

# Пример: парсим run_results.json и отправляем метрики в Prometheus
import json
import time
from prometheus_client import Gauge, push_to_gateway

def push_dbt_metrics(run_results_path):
    with open(run_results_path) as f:
        results = json.load(f)

    model_duration = Gauge(
        'dbt_model_duration_seconds',
        'dbt model execution time',
        ['model_name', 'status']
    )

    for result in results['results']:
        model_duration.labels(
            model_name=result['unique_id'],
            status=result['status']
        ).set(result['execution_time'])

    push_to_gateway('prometheus-pushgateway:9091', job='dbt', registry=registry)

Лабораторная практика: сквозной CI/CD пайплайн

Бизнес-кейс

Настроим автоматический ежедневный пересчёт маркетинговой витрины с контролем качества. Требования:

  1. Витрина fct_marketing_events пересчитывается каждый день в 06:00 UTC
  2. Перед пересчётом - проверка свежести источников (warn: 1 час, error: 4 часа)
  3. Любой PR должен пройти линтинг, компиляцию и тест изменённых моделей
  4. Merge в main автоматически деплоит новую версию на Airflow-воркеры

Шаг 1: Создание тегов в dbt_project.yml

# dbt_project.yml
models:
  my_project:
    staging:
      +tags: ['staging']
      +materialized: view

    intermediate:
      +tags: ['intermediate']
      +materialized: ephemeral

    marts:
      +tags: ['marts']
      marketing:
        +tags: ['marketing']
        +materialized: incremental
        +file_format: parquet
        +incremental_strategy: insert_overwrite
        +location_root: "{{ env_var('DBT_LOCATION_ROOT') }}"

Шаг 2: DAG в Airflow

# dags/dbt_marketing_pipeline.py
from datetime import datetime, timedelta
from airflow import DAG
from airflow.operators.bash import BashOperator
from airflow.hooks.base import BaseHook
from airflow.models import Variable

def get_dbt_env():
    conn = BaseHook.get_connection("spark_thrift_prod")
    return {
        "DBT_TARGET": "prod",
        "SPARK_HOST": conn.host,
        "SPARK_PORT": str(conn.port or 10000),
        "SPARK_USER": conn.login,
        "SPARK_PASSWORD": conn.password,
        "DBT_LOCATION_ROOT": Variable.get("dbt_location_root_prod"),
    }

with DAG(
    dag_id="dbt_marketing_pipeline",
    description="Ежедневный пересчёт маркетинговых витрин",
    schedule="0 6 * * *",
    start_date=datetime(2024, 1, 1),
    catchup=False,
    max_active_runs=1,
    default_args={
        "owner": "marketing-data",
        "retries": 2,
        "retry_delay": timedelta(minutes=10),
        "email_on_failure": True,
        "email": ["marketing-data@company.ru"],
        "sla": timedelta(hours=2),
    },
    tags=["dbt", "marketing", "daily"],
) as dag:

    check_freshness = BashOperator(
        task_id="check_source_freshness",
        bash_command=(
            "cd /opt/dbt && "
            "dbt source freshness "
            "--profiles-dir /opt/dbt "
            "--select 'source:bronze_marketing'"
        ),
        env=get_dbt_env(),
    )

    build_staging = BashOperator(
        task_id="build_staging_layer",
        bash_command=(
            "cd /opt/dbt && "
            "dbt build "
            "--profiles-dir /opt/dbt "
            "--select 'tag:staging'"
        ),
        env=get_dbt_env(),
    )

    build_marketing = BashOperator(
        task_id="build_marketing_marts",
        bash_command=(
            "cd /opt/dbt && "
            "dbt build "
            "--profiles-dir /opt/dbt "
            "--select 'tag:marketing'"
        ),
        env=get_dbt_env(),
    )

    check_freshness >> build_staging >> build_marketing

Шаг 3: CI workflow

# .github/workflows/ci.yml
name: dbt CI

on:
  pull_request:
    branches: [main]
    paths: ['models/**', 'tests/**', 'macros/**', 'dbt_project.yml']

jobs:
  validate:
    runs-on: ubuntu-latest
    steps:
      - uses: actions/checkout@v4

      - name: Setup Python
        uses: actions/setup-python@v5
        with:
          python-version: '3.11'
          cache: 'pip'

      - name: Install dependencies
        run: |
          pip install dbt-spark[PyHive]==1.7.4 sqlfluff==2.3.0 sqlfluff-templater-dbt==2.3.0
          dbt deps

      - name: SQL Lint
        run: sqlfluff lint models/ --dialect sparksql --templater dbt

      - name: dbt compile
        run: dbt compile --target ci
        env:
          DBT_TARGET: ci
          SPARK_HOST: ${{ secrets.SPARK_CI_HOST }}
          SPARK_PASSWORD: ${{ secrets.SPARK_CI_PASSWORD }}

      - name: Download prod manifest
        run: aws s3 cp s3://dbt-artifacts/prod/manifest.json ./prod_manifest/manifest.json
        env:
          AWS_ACCESS_KEY_ID: ${{ secrets.AWS_ACCESS_KEY_ID }}
          AWS_SECRET_ACCESS_KEY: ${{ secrets.AWS_SECRET_ACCESS_KEY }}

      - name: dbt build (modified only)
        run: |
          dbt build \
            --target ci \
            --select state:modified+ \
            --defer \
            --state ./prod_manifest/
        env:
          DBT_TARGET: ci
          SPARK_HOST: ${{ secrets.SPARK_CI_HOST }}
          SPARK_PASSWORD: ${{ secrets.SPARK_CI_PASSWORD }}
          DBT_LOCATION_ROOT: s3a://ci-lake/gold
          PR_NUMBER: ${{ github.event.number }}

Шаг 4: Симуляция падения CI при нарушении Data Contract

Представим, что разработчик меняет тип колонки в модели:

-- БЫЛО:
CAST(amount / 100.0 AS DECIMAL(18, 2)) AS amount_rub

-- СТАЛО (разработчик убрал CAST):
amount / 100.0 AS amount_rub  -- тип станет DOUBLE вместо DECIMAL(18,2)

При следующем dbt build в CI:

Compilation Error in model fct_marketing_events
  This model has an enforced contract that failed.
  Column "amount_rub": data type mismatch
    expected: decimal(18,2)
    got:      double

  Please update your model to ensure the data types match.

16:45:23  Encountered an error:
  Compilation failed

❌ CI FAILED - PR blocked from merging

GitHub покажет красный крестик на PR. Разработчик возвращает CAST:

CAST(amount / 100.0 AS DECIMAL(18, 2)) AS amount_rub

Следующий CI-ран проходит, PR получает зелёную галочку и может быть merged.

Шаг 5: Настройка защиты ветки в GitHub

В настройках репозитория (Settings → Branches → Add rule):

  • Branch name pattern: main
  • Require status checks to pass before merging: ✓
  • Required status checks: validate (из CI workflow)
  • Require up-to-date branches: ✓
  • Do not allow bypassing the above settings: ✓

Теперь ни один PR не может быть merged в main без прохождения CI. Защита production от некорректного кода.


Планирование и scheduling-стратегии

Cron-based scheduling

Стандартное расписание для batch-пайплайнов:

schedule="0 6 * * *"        # Каждый день в 06:00 UTC
schedule="0 */4 * * *"      # Каждые 4 часа
schedule="30 5 * * 1"       # Каждый понедельник в 05:30 UTC (еженедельный отчёт)
schedule="0 0 1 * *"        # Первое число каждого месяца

Data-aware scheduling (Airflow 2.4+)

Airflow поддерживает Datasets - триггер DAG по событию обновления данных:

from airflow.datasets import Dataset

# Определяем датасет
TRANSACTIONS_DATASET = Dataset("s3://bronze/raw_transactions/")

# DAG, который обновляет датасет
with DAG("ingestion_pipeline", schedule="*/15 * * * *") as ingest_dag:
    ingest = BashOperator(
        task_id="ingest",
        bash_command="python /apps/ingest_transactions.py",
        outlets=[TRANSACTIONS_DATASET],  # Помечаем: этот таск обновляет датасет
    )

# DAG, который запускается при обновлении датасета
with DAG(
    "dbt_streaming_pipeline",
    schedule=[TRANSACTIONS_DATASET],  # Запускаться при обновлении датасета
) as dbt_dag:
    build = BashOperator(
        task_id="build_marts",
        bash_command="dbt build --select tag:finance",
    )

С Data-aware scheduling dbt запускается не по расписанию, а по факту появления новых данных. Это устраняет искусственные задержки: не ждём 06:00 если данные пришли в 05:15.

Идемпотентность и повторные запуски

В distributed системах всё может упасть и перезапуститься. Airflow поддерживает автоматические retry. Чтобы retry работал корректно - таски и dbt-модели должны быть идемпотентными: повторный запуск даёт тот же результат, что и первый.

dbt с incremental_strategy='insert_overwrite' идемпотентен: повторное выполнение просто перезапишет партицию теми же данными. С append - не идемпотентен: повторное выполнение создаёт дубли.

# Правильная настройка retry:
default_args={
    "retries": 3,
    "retry_delay": timedelta(minutes=5),
    "retry_exponential_backoff": True,  # 5min → 10min → 20min
    "max_retry_delay": timedelta(hours=1),
}

Паттерны и антипаттерны

Антипаттерн 1: Монолитный DAG с одним таском

# ПЛОХО: весь dbt в одном таске - непрозрачно, нельзя перезапустить часть
run_all_dbt = BashOperator(
    bash_command="dbt build --select '+'"
)

Если упала одна модель из 200 - нужно перезапускать всё. Нет видимости, что именно упало.

Лучше: разбить на логические группы по слоям или доменам.

Антипаттерн 2: Пропуск тестов

# ПЛОХО: run без тестов → некачественные данные попадают в Gold
run_models = BashOperator(bash_command="dbt run --select '+fct_orders'")
# Тесты вообще не запускаются!

Лучше: всегда dbt build в production.

Антипаттерн 3: Hardcode credentials в DAG

# ПЛОХО: пароль прямо в коде
env={
    "SPARK_PASSWORD": "super_secret_password_123",
}

Лучше: Airflow Connections + Secrets Backend.

Антипаттерн 4: Отсутствие max_active_runs=1

Без этого параметра Airflow может запустить несколько ранов одного DAG параллельно (catchup или ручной триггер). Два паралельных dbt build на одни и те же партиции - race condition, перезапись данных в процессе чтения.

with DAG(
    ...,
    max_active_runs=1,   # Всегда ставить для dbt DAG
) as dag:
    ...

Антипаттерн 5: CI без Slim CI

# ПЛОХО: каждый PR пересчитывает весь проект
- run: dbt build --target ci

# ХОРОШО: только изменённые модели
- run: dbt build --target ci --select state:modified+ --defer --state ./prod_manifest/

Без Slim CI CI-ран стоит столько же, сколько production-ран. С Slim CI - в 10-100 раз дешевле.


Домашнее задание

Задача 1: BashOperator DAG

Создайте Airflow DAG dbt_finance_pipeline с последовательными шагами:

  1. check_freshness - dbt source freshness --select source:bronze
  2. build_staging - dbt build --select tag:staging
  3. build_intermediate - dbt build --select tag:intermediate
  4. build_finance и build_marketing - параллельно (оба зависят от build_intermediate)

Требования:

  • Расписание: каждый день в 07:00 UTC
  • max_active_runs=1
  • Retry: 2 попытки с задержкой 10 минут
  • Credentials через Airflow Connections (не хардкод)

Задача 2: Credentials через Airflow Connection

Настройте Airflow Connection spark_thrift_prod:

  • Тип: Generic / HTTP
  • Host: адрес Thrift Server
  • Login: airflow_svc
  • Password: сервисный пароль

Обновите DAG так, чтобы credentials читались через BaseHook.get_connection("spark_thrift_prod") и передавались в dbt как переменные окружения.

Задача 3: GitHub Actions CI

Напишите .github/workflows/ci.yml для вашего dbt-проекта:

  1. Триггер: Pull Request в main
  2. Шаги: установка dbt → dbt compile → скачать prod manifest → dbt build --select state:modified+ --defer
  3. Секреты: SPARK_CI_HOST и SPARK_CI_PASSWORD из GitHub Secrets

Задача 4: Защита ветки

Опишите текстово (или настройте в реальном репозитории) правила защиты ветки main:

  1. Обязательные статус-чеки: CI workflow из задачи 3
  2. Требуется минимум 1 review
  3. Запрещена прямая пуш в main (только через PR)

Объясните, как эта защита предотвращает попадание моделей с нарушенным unique-ключом в production.

Что сдавать

  1. Код dags/dbt_finance_pipeline.py
  2. Файл .github/workflows/ci.yml
  3. Скриншот Grid View в Airflow UI с успешно отработавшим DAG
  4. Скриншот GitHub PR со статусом CI (зелёный ✓ или красный ✗ с описанием)

Задача со звёздочкой: внедрите Slim CI (--defer --state) и покажите разницу во времени выполнения CI с и без него на проекте из 10+ моделей.