dbt + Airflow: BashOperator dbt run vs SparkSubmitOperator, CI/CD pipeline
dbt + Airflow: BashOperator dbt run vs SparkSubmitOperator, CI/CD pipeline
Архитектура оркестрации 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.
Решения:
- Не парсить manifest.json в DAG-файле - задать модели вручную или группами (по тегам)
- Использовать Cosmos - специализированная библиотека, которая делает это эффективно
- Использовать 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 пайплайн¶
Бизнес-кейс¶
Настроим автоматический ежедневный пересчёт маркетинговой витрины с контролем качества. Требования:
- Витрина
fct_marketing_eventsпересчитывается каждый день в 06:00 UTC - Перед пересчётом - проверка свежести источников (warn: 1 час, error: 4 часа)
- Любой PR должен пройти линтинг, компиляцию и тест изменённых моделей
- 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 с последовательными шагами:
check_freshness-dbt source freshness --select source:bronzebuild_staging-dbt build --select tag:stagingbuild_intermediate-dbt build --select tag:intermediatebuild_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-проекта:
- Триггер: Pull Request в
main - Шаги: установка dbt →
dbt compile→ скачать prod manifest →dbt build --select state:modified+ --defer - Секреты:
SPARK_CI_HOSTиSPARK_CI_PASSWORDиз GitHub Secrets
Задача 4: Защита ветки¶
Опишите текстово (или настройте в реальном репозитории) правила защиты ветки main:
- Обязательные статус-чеки: CI workflow из задачи 3
- Требуется минимум 1 review
- Запрещена прямая пуш в
main(только через PR)
Объясните, как эта защита предотвращает попадание моделей с нарушенным unique-ключом в production.
Что сдавать¶
- Код
dags/dbt_finance_pipeline.py - Файл
.github/workflows/ci.yml - Скриншот Grid View в Airflow UI с успешно отработавшим DAG
- Скриншот GitHub PR со статусом CI (зелёный ✓ или красный ✗ с описанием)
Задача со звёздочкой: внедрите Slim CI (--defer --state) и покажите разницу во времени выполнения CI с и без него на проекте из 10+ моделей.