dbt Core: models, refs, sources и materialization types (table, view, incremental)
dbt Core превращает SELECT-запросы в управляемый DAG трансформаций поверх Spark SQL. Разбираем sources, ref(), стратегии материализации и incremental-модели для lakehouse.
Зачем dbt в мире Big Data и Spark¶
Проблема императивного кода¶
До появления dbt типичный аналитический инженер писал трансформации данных в виде Python-скриптов с PySpark или SQL-скриптов с явными DDL-командами. Пайплайн для расчёта дневных метрик пользователей выглядел примерно так:
from pyspark.sql import SparkSession
from pyspark.sql.functions import col, sum as spark_sum, to_date, count
spark = SparkSession.builder.appName("User Metrics ETL").getOrCreate()
raw_events = spark.read.parquet("s3a://bronze/events/")
clean_events = raw_events \
.filter(col("event_type").isNotNull()) \
.filter(col("user_id") > 0) \
.withColumn("event_date", to_date(col("event_ts"))) \
.select("user_id", "session_id", "event_type", "event_date", "revenue")
clean_events.createOrReplaceTempView("stg_events")
user_metrics = spark.sql("""
SELECT
user_id,
event_date,
COUNT(DISTINCT session_id) AS sessions,
SUM(revenue) AS total_revenue,
COUNT(*) AS event_count
FROM stg_events
GROUP BY user_id, event_date
""")
user_metrics.write \
.format("delta") \
.mode("overwrite") \
.partitionBy("event_date") \
.saveAsTable("gold.fct_user_activity")
На первый взгляд всё понятно. Но в production с десятками таких скриптов начинаются проблемы:
Проблема зависимостей: скрипты зависят друг от друга - fct_user_activity нужен stg_events, который нужен raw_events. Но эти зависимости нигде явно не зафиксированы. Если кто-то запустит скрипты в неправильном порядке, всё сломается без понятного сообщения об ошибке.
Проблема хардкода: пути s3a://bronze/events/, имя таблицы gold.fct_user_activity, имя базы данных gold - всё это зашито в код. На dev-окружении нужно переключаться на dev_gold.fct_user_activity, в staging - на staging_gold. Каждое изменение окружения требует ручной правки кода или сложной параметризации.
Проблема воспроизводимости: нет единого места, где видно "что из чего получается". Линейность графа зависимостей нужно восстанавливать вручную, читая каждый скрипт.
Проблема тестируемости: чтобы убедиться, что user_id никогда не NULL в выходной таблице, нужно написать отдельный скрипт для проверки. Или просто надеяться, что всё хорошо.
Декларативный подход dbt¶
dbt (data build tool) предлагает принципиально иной подход: аналитик пишет только SELECT-запрос, описывающий бизнес-логику, а dbt берёт на себя всё остальное - создание таблиц, управление зависимостями, тестирование, документирование.
-- models/gold/fct_user_activity.sql
{{ config(
materialized='incremental',
unique_key=['user_id', 'event_date'],
partition_by={'field': 'event_date', 'data_type': 'date'}
) }}
SELECT
user_id,
event_date,
COUNT(DISTINCT session_id) AS sessions,
SUM(revenue) AS total_revenue,
COUNT(*) AS event_count
FROM {{ ref('stg_events') }}
{% if is_incremental() %}
WHERE event_date >= (SELECT MAX(event_date) FROM {{ this }})
{% endif %}
GROUP BY user_id, event_date
Это тот же запрос, но теперь:
{{ ref('stg_events') }}- явная зависимость, dbt автоматически построит граф и выполнитstg_eventsпервым{{ config(...) }}- инфраструктурные параметры (формат хранения, партиционирование) отделены от логики{% if is_incremental() %}- логика инкрементального обновления встроена прямо в SQL- Нет ни одного DDL-оператора, нет хардкода путей, нет зависимости от конкретного окружения
Место dbt в архитектуре¶
Диаграмма показывает чёткое разделение ответственности: Ingestion (загрузка в Bronze) - это задача Airflow, Dagster, Airbyte, Spark Streaming. Они умеют читать API, слушать Kafka, делать CDC из баз данных. Transformation (Silver и Gold) - это зона dbt. dbt работает только с данными, которые уже лежат в хранилище, превращая их в аналитически пригодный вид через SQL. Serving (отдача в BI) - это задача Gold-слоя.
Такое разделение называется ELT (Extract-Load-Transform) в отличие от старого ETL (Extract-Transform-Load). В ELT данные сначала загружаются "как есть" (Extract + Load), а затем трансформируются уже внутри хранилища (Transform). dbt реализует шаг Transform в этой схеме.
Адаптер dbt-spark¶
dbt Core - это фреймворк с поддержкой разных "адаптеров" для разных SQL-движков. Для работы со Spark существует адаптер dbt-spark, который поддерживает несколько режимов подключения:
| Режим подключения | Описание | Когда использовать |
|---|---|---|
thrift |
Spark Thrift Server (JDBC) | Self-hosted Spark кластер |
http |
Spark Connect / Databricks SQL | Databricks, cloud Spark |
session |
Embedded SparkSession | Тестирование локально |
pip install dbt-spark[PyHive]
dbt --version
Архитектура dbt-проекта¶
Структура директорий¶
my_dbt_project/
├── dbt_project.yml # Главный конфиг проекта
├── profiles.yml # Подключения к БД (обычно в ~/.dbt/)
├── models/ # SQL-трансформации
│ ├── staging/ # stg_* - нормализация Bronze
│ │ ├── sources.yml # Описание источников
│ │ ├── stg_events.sql
│ │ └── stg_users.sql
│ ├── intermediate/ # int_* - промежуточная логика
│ │ └── int_sessions.sql
│ └── marts/ # fct_*, dim_* - Gold-витрины
│ ├── fct_user_activity.sql
│ └── dim_campaigns.sql
├── macros/ # Jinja-макросы
│ └── generate_schema.sql
├── tests/ # Кастомные тесты
│ └── assert_revenue_positive.sql
├── seeds/ # CSV-файлы (справочники)
│ └── campaign_types.csv
└── snapshots/ # SCD Type 2 (история изменений)
└── snap_user_status.sql
Каждая директория - это слой Medallion-архитектуры:
- staging/ - модели с префиксом
stg_. Работают с источниками напрямую: переименовывают колонки, кастируют типы, убирают дубли. Один источник - одна staging-модель. - intermediate/ - модели с префиксом
int_. Реализуют бизнес-логику, которая используется несколькими marts. Не предназначены для конечных пользователей. - marts/ - модели с префиксом
fct_(факты) илиdim_(измерения). Готовы для BI-инструментов. Оптимизированы для чтения.
dbt_project.yml - главный конфиг¶
name: 'analytics_lakehouse'
version: '1.0.0'
config-version: 2
profile: 'spark_local'
model-paths: ["models"]
test-paths: ["tests"]
macro-paths: ["macros"]
seed-paths: ["seeds"]
snapshot-paths: ["snapshots"]
models:
analytics_lakehouse:
staging:
+materialized: view
+schema: staging
intermediate:
+materialized: ephemeral
+schema: intermediate
marts:
+materialized: table
+schema: marts
fct_user_activity:
+materialized: incremental
Блок models: - это конфигурация по умолчанию для каждой поддиректории. +materialized: view для staging означает: все модели в папке staging/ по умолчанию создаются как views, если не указано иное в самой модели. Знак + - это "применить к директории и всем вложенным".
profiles.yml - подключение к Spark¶
spark_local:
target: dev
outputs:
dev:
type: spark
method: thrift
host: localhost
port: 10000
user: hadoop
schema: dbt_dev
connect_retries: 3
connect_timeout: 10
server_side_parameters:
spark.sql.shuffle.partitions: "200"
spark.sql.adaptive.enabled: "true"
prod:
type: spark
method: thrift
host: spark-thrift.internal
port: 10000
user: dbt_service
schema: dbt_prod
server_side_parameters:
spark.sql.shuffle.partitions: "2000"
spark.sql.adaptive.enabled: "true"
spark.sql.adaptive.advisoryPartitionSizeInBytes: "134217728"
schema в контексте Spark - это имя базы данных (database/namespace). В dev-окружении все таблицы создаются в базе dbt_dev, в prod - в dbt_prod. Это обеспечивает изоляцию: разработчик не может случайно перезаписать production-таблицы.
server_side_parameters позволяет передавать Spark-конфигурации через JDBC-соединение. Это удобно для установки разных значений shuffle.partitions для dev и prod без изменения кода.
Sources: описание источников¶
Зачем нужны Sources¶
Sources в dbt - это механизм регистрации "сырых" таблиц или файлов, которые dbt не создаёт, а только читает. Это Bronze-слой: данные, загруженные Airflow/Airbyte/Spark Streaming.
Без Sources аналитик написал бы в модели:
-- stg_events.sql - ПЛОХО
SELECT * FROM bronze.raw_events WHERE event_date = '2024-01-15'
Проблемы такого подхода:
bronze.raw_events- хардкод. На dev может бытьdev_bronze.raw_events, в другом кластере - другой путь.- Нет возможности проверить freshness: dbt не знает, что
bronze.raw_events- это внешний источник. - Lineage-граф не покажет эту зависимость - dbt не "видит" таблицы без регистрации.
Описание Sources в sources.yml¶
# models/staging/sources.yml
version: 2
sources:
- name: bronze
description: "Bronze-слой: сырые данные от ingestion-пайплайнов"
database: bronze_db
schema: events
tables:
- name: raw_events
description: "Сырые события пользователей из Kafka"
identifier: raw_events
loaded_at_field: event_ts
freshness:
warn_after: {count: 1, period: hour}
error_after: {count: 3, period: hour}
columns:
- name: event_id
description: "Уникальный ID события"
tests:
- not_null
- unique
- name: user_id
description: "ID пользователя"
tests:
- not_null
- name: event_type
description: "Тип события"
- name: event_ts
description: "Временная метка события"
tests:
- not_null
- name: revenue
description: "Выручка (только для purchase событий)"
- name: raw_users
description: "Справочник пользователей из PostgreSQL CDC"
identifier: raw_users
loaded_at_field: updated_at
freshness:
warn_after: {count: 6, period: hour}
error_after: {count: 24, period: hour}
Разберём ключевые поля:
name - логическое имя источника в dbt. Используется в {{ source('bronze', 'raw_events') }}.
database и schema - указывают на реальное местоположение таблицы в Spark Metastore. Для Spark это database.table.
identifier - реальное имя таблицы, если оно отличается от name. Позволяет дать короткое логическое имя таблице с длинным техническим именем.
loaded_at_field - колонка с временной меткой последней загрузки. Используется командой dbt source freshness для проверки актуальности данных.
freshness - пороги свежести. Если данные не обновлялись более 1 часа - warn, более 3 часов - error. Это первая линия защиты от ситуации "забыли запустить ingestion-пайплайн, а BI уже использует устаревшие данные".
Использование source() в модели¶
-- models/staging/stg_events.sql
{{ config(materialized='view') }}
SELECT
event_id,
user_id,
session_id,
event_type,
CAST(event_ts AS TIMESTAMP) AS event_ts,
CAST(event_ts AS DATE) AS event_date,
COALESCE(revenue, 0.0) AS revenue,
LOWER(TRIM(device_type)) AS device_type,
CURRENT_TIMESTAMP() AS dbt_loaded_at
FROM {{ source('bronze', 'raw_events') }}
WHERE event_id IS NOT NULL
AND user_id IS NOT NULL
{{ source('bronze', 'raw_events') }} - это Jinja-макрос. При компиляции dbt подставит сюда полное имя таблицы: bronze_db.events.raw_events. На dev-профиле это может быть dev_bronze_db.events.raw_events - в зависимости от конфигурации.
Проверить freshness всех источников:
dbt source freshness
Вывод:
11:23:45 Running with dbt=1.7.0
11:23:45 Found 3 sources, 0 freshness errors, 0 freshness warnings
11:23:46 bronze.raw_events ... [PASS (23 minutes ago)]
11:23:46 bronze.raw_users ... [WARN (4 hours ago)]
Предупреждение для raw_users - данные не обновлялись 4 часа, порог предупреждения 6 часов. Следует проверить ingestion-пайплайн до того, как данные устареют критически.
Models: SQL как единица трансформации¶
Модель - это SELECT-запрос¶
Ключевая идея dbt: модель - это просто SELECT-запрос. Никакого DDL (CREATE TABLE, INSERT INTO, CREATE VIEW). Аналитик описывает, что он хочет получить, а dbt решает как это создать в базе данных.
-- models/staging/stg_users.sql
{{ config(materialized='view') }}
SELECT
user_id,
LOWER(email) AS email,
COALESCE(first_name, 'Unknown') AS first_name,
COALESCE(last_name, 'Unknown') AS last_name,
CAST(created_at AS TIMESTAMP) AS created_at,
CAST(created_at AS DATE) AS signup_date,
CASE
WHEN account_type = 'premium' THEN TRUE
WHEN account_type = 'trial' THEN TRUE
ELSE FALSE
END AS is_paid_user,
CURRENT_TIMESTAMP() AS dbt_loaded_at
FROM {{ source('bronze', 'raw_users') }}
WHERE user_id IS NOT NULL
AND email IS NOT NULL
Этот файл - и есть вся модель stg_users. Когда dbt выполнит dbt run, он:
- Скомпилирует Jinja-шаблоны, подставив реальные имена таблиц
- Создаст VIEW (потому что
materialized='view') в базе данных:CREATE OR REPLACE VIEW dbt_dev.stg_users AS SELECT ... - Зарегистрирует зависимость:
stg_usersзависит отsource('bronze', 'raw_users')
Jinja templating: динамический SQL¶
dbt использует Jinja - шаблонизатор для Python. В SQL-файлах можно использовать переменные, условия, циклы и макросы.
{{ config(...) }}
SELECT ...
{% if condition %}
WHERE ...
{% endif %}
Двойные фигурные скобки {{ }} - вставка выражения. Блок {% %} - управляющая конструкция (if, for, set).
Встроенные переменные dbt:
SELECT
'{{ env_var("DBT_ENV", "dev") }}' AS environment,
'{{ var("run_date", "2024-01-01") }}' AS run_date,
CURRENT_TIMESTAMP() AS compiled_at
FROM ...
env_var() читает переменную окружения; var() читает переменную, переданную в командной строке: dbt run --vars '{"run_date": "2024-01-15"}'.
ref(): управление зависимостями¶
Что делает ref()¶
{{ ref('model_name') }} - фундаментальный макрос dbt. Он выполняет две задачи одновременно:
1. Подстановка правильного имени таблицы: на dev-профиле {{ ref('stg_events') }} превратится в dbt_dev.stg_events, на prod - в dbt_prod.stg_events. Никакого хардкода окружения.
2. Регистрация зависимости: dbt видит каждый ref() и строит граф зависимостей. Если модель fct_user_activity содержит {{ ref('stg_events') }}, dbt знает: сначала нужно выполнить stg_events, потом fct_user_activity.
-- models/intermediate/int_sessions.sql
{{ config(materialized='ephemeral') }}
SELECT
user_id,
session_id,
MIN(event_ts) AS session_start,
MAX(event_ts) AS session_end,
MAX(event_ts) - MIN(event_ts) AS session_duration,
COUNT(*) AS event_count,
SUM(revenue) AS session_revenue,
COUNT(DISTINCT event_type) AS unique_event_types,
CAST(MIN(event_ts) AS DATE) AS session_date
FROM {{ ref('stg_events') }}
GROUP BY user_id, session_id
-- models/marts/fct_user_activity.sql
{{ config(materialized='table') }}
SELECT
e.user_id,
u.email,
u.is_paid_user,
e.event_date,
COUNT(DISTINCT e.session_id) AS sessions,
SUM(s.session_duration) AS total_session_time,
SUM(e.revenue) AS total_revenue,
COUNT(*) AS event_count
FROM {{ ref('stg_events') }} AS e
JOIN {{ ref('stg_users') }} AS u USING (user_id)
JOIN {{ ref('int_sessions') }} AS s USING (user_id, session_id)
GROUP BY e.user_id, u.email, u.is_paid_user, e.event_date
В этом примере fct_user_activity ссылается на три модели. dbt автоматически построит граф:
sources/bronze/raw_events → stg_events ─┬─→ int_sessions ─┐
│ ↓
sources/bronze/raw_users → stg_users ──┴────────────→ fct_user_activity
И выполнит их в правильном порядке, с возможностью параллельного выполнения независимых веток.
Просмотр скомпилированного SQL¶
После {{ ref() }} в скомпилированном SQL появятся реальные имена таблиц. Чтобы посмотреть, что именно выполняется:
dbt compile --select fct_user_activity
Скомпилированный SQL сохраняется в target/compiled/:
-- target/compiled/analytics_lakehouse/models/marts/fct_user_activity.sql
SELECT
e.user_id,
u.email,
...
FROM dbt_dev.stg_events AS e
JOIN dbt_dev.stg_users AS u USING (user_id)
JOIN dbt_dev.int_sessions AS s USING (user_id, session_id)
GROUP BY e.user_id, u.email, u.is_paid_user, e.event_date
Это полезно для отладки: можно скопировать скомпилированный SQL и запустить его напрямую в Spark SQL или Beeline.
Data Lineage: визуализация зависимостей¶
dbt docs generate
dbt docs serve
После этих команд в браузере откроется интерактивная документация со встроенным Lineage Graph - визуальным графом зависимостей. Каждый узел - модель или источник. Стрелки показывают направление данных. Можно кликнуть на модель и увидеть её описание, колонки, тесты и скомпилированный SQL.
Lineage Graph - это не просто документация. Это инструмент анализа влияния изменений: если нужно изменить схему Bronze-таблицы, через граф сразу видно, какие downstream-модели затронуты.
Стратегии материализации¶
Материализация (materialization) определяет, как dbt физически создаёт результат модели в базе данных. Это один из самых важных архитектурных решений в dbt.
Диаграмма показывает четыре типа материализации. Каждый тип - это разный способ, которым dbt создаёт объект в базе данных. Выбор типа определяется объёмом данных, частотой обновления и требованиями к скорости чтения.
View: логическая абстракция¶
Материализация view - самая лёгкая. dbt создаёт обычный SQL VIEW:
CREATE OR REPLACE VIEW dbt_dev.stg_events AS
SELECT
event_id,
user_id,
...
FROM bronze_db.events.raw_events
WHERE event_id IS NOT NULL
Преимущества:
- Нет дополнительного хранилища - данные не копируются
- Всегда актуальны - каждый запрос к view читает свежие данные из источника
- Мгновенное создание -
dbt runвыполняется за секунды
Ограничения:
- Каждый запрос к view выполняет всю трансформацию заново
- Если view используется в нескольких downstream-моделях, трансформация повторяется для каждой из них
- Нельзя партиционировать - view не имеет физического хранилища
Когда использовать:
- Staging-модели: они - тонкая обёртка над Bronze, используются один раз на следующем уровне
- Лёгкие фильтры без тяжёлых агрегаций
- Когда актуальность данных критичнее производительности
-- models/staging/stg_events.sql
{{ config(materialized='view') }}
SELECT
event_id,
user_id,
CAST(event_ts AS TIMESTAMP) AS event_ts,
CAST(event_ts AS DATE) AS event_date,
event_type,
COALESCE(revenue, 0.0) AS revenue
FROM {{ source('bronze', 'raw_events') }}
WHERE event_id IS NOT NULL
Ephemeral: CTE в памяти¶
ephemeral - особый тип: dbt не создаёт никакого объекта в базе данных. Вместо этого при компиляции SQL код модели вставляется как CTE (Common Table Expression) прямо в тело зависящей модели.
-- models/intermediate/int_sessions.sql
{{ config(materialized='ephemeral') }}
SELECT
user_id,
session_id,
MIN(event_ts) AS session_start,
MAX(event_ts) AS session_end,
COUNT(*) AS event_count
FROM {{ ref('stg_events') }}
GROUP BY user_id, session_id
Когда fct_user_activity ссылается на int_sessions через {{ ref('int_sessions') }}, скомпилированный SQL будет содержать CTE:
-- Скомпилированный fct_user_activity.sql
WITH int_sessions AS (
SELECT
user_id,
session_id,
MIN(event_ts) AS session_start,
MAX(event_ts) AS session_end,
COUNT(*) AS event_count
FROM dbt_dev.stg_events
GROUP BY user_id, session_id
)
SELECT
e.user_id,
...
FROM dbt_dev.stg_events AS e
JOIN int_sessions AS s USING (user_id, session_id)
Преимущества:
- Нет лишних объектов в базе данных
- Промежуточная логика остаётся инкапсулированной
Ограничения:
- Логика дублируется, если на ephemeral-модель ссылаются несколько моделей
- Нельзя запустить
dbt run --select int_sessions- не создаёт ничего напрямую - Нет возможности отладить - результат нигде не сохраняется
Когда использовать: для промежуточных шагов, которые используются только одной downstream-моделью и не нужны независимо.
Table: полный пересчёт¶
Материализация table - полный пересчёт таблицы при каждом запуске:
-- Что генерирует dbt при table materialization
CREATE TABLE dbt_dev.fct_revenue_summary AS
SELECT
campaign_id,
channel,
event_date,
SUM(revenue) AS total_revenue,
COUNT(*) AS conversions
FROM dbt_dev.stg_events
GROUP BY campaign_id, channel, event_date
В Spark dbt использует паттерн CTAS + RENAME или CREATE TABLE AS SELECT. На практике это означает:
- dbt создаёт временную таблицу
dbt_dev.fct_revenue_summary__dbt_tmp - Выполняет SELECT в эту временную таблицу
- Переименовывает временную таблицу в
dbt_dev.fct_revenue_summary(атомарная операция)
Такой подход обеспечивает атомарность: читатели видят либо старую версию таблицы, либо новую, но никогда - частичную.
Преимущества:
- Простота: нет инкрементальной логики, нет состояния
- Быстрое чтение: данные физически сохранены и оптимизированы
- Можно задать партиционирование и формат хранения
Ограничения:
- Полный пересчёт при каждом запуске: на терабайтных данных это часы работы кластера
- Нет возможности обработать только "сегодняшние" данные
Настройка для Spark:
{{ config(
materialized='table',
file_format='delta',
partition_by=['event_date'],
clustered_by=['campaign_id'],
buckets=8,
location='s3a://datalake/gold/fct_revenue_summary'
) }}
SELECT
campaign_id,
channel,
event_date,
SUM(revenue) AS total_revenue,
COUNT(*) AS conversions
FROM {{ ref('stg_events') }}
GROUP BY campaign_id, channel, event_date
file_format='delta' создаст Delta-таблицу вместо Parquet. partition_by - Spark PARTITIONED BY clause. location - путь в объектном хранилище.
Когда использовать:
- Gold-витрины с умеренным объёмом данных (до нескольких ГБ после агрегации)
- Справочные (dimension) таблицы
- Когда нужна простота без incremental-логики
Incremental: обработка только новых данных¶
Incremental - самая мощная и сложная стратегия. Идея: при первом запуске таблица создаётся полностью. При последующих запусках обрабатываются только новые данные и дозаписываются (или обновляются) в существующую таблицу.
Это критически важно для больших данных: вместо полного пересчёта терабайтного датасета обрабатывается только прирост за последние сутки - в 365 раз меньше вычислений.
Базовая структура incremental-модели¶
{{ config(
materialized='incremental',
unique_key='event_id',
incremental_strategy='merge',
file_format='delta',
partition_by=['event_date'],
) }}
SELECT
event_id,
user_id,
session_id,
event_type,
event_date,
revenue,
CURRENT_TIMESTAMP() AS processed_at
FROM {{ ref('stg_events') }}
{% if is_incremental() %}
WHERE event_date >= (SELECT MAX(event_date) FROM {{ this }})
{% endif %}
Разберём каждую часть:
{% if is_incremental() %} - блок выполняется только при инкрементальном запуске (не при первом создании таблицы и не при --full-refresh). На первом запуске этот блок игнорируется, и модель обрабатывает все данные.
{{ this }} - ссылка на саму текущую таблицу. На первом запуске таблицы не существует, поэтому блок is_incremental() не выполняется. После первого запуска {{ this }} вернёт реальное имя таблицы в базе.
SELECT MAX(event_date) FROM {{ this }} - watermark: мы читаем максимальную дату, уже имеющуюся в таблице, и обрабатываем только данные с этой даты и позже. Перекрытие (включение текущей max-даты) сделано намеренно: события за сегодня могут ещё поступать.
Стратегии incremental в dbt-spark¶
Параметр incremental_strategy определяет механизм добавления данных. dbt-spark поддерживает три стратегии:
append - простая дозапись
{{ config(
materialized='incremental',
incremental_strategy='append',
partition_by=['event_date'],
file_format='parquet',
) }}
SELECT
event_id,
user_id,
event_date,
revenue
FROM {{ ref('stg_events') }}
{% if is_incremental() %}
WHERE event_date > (SELECT MAX(event_date) FROM {{ this }})
{% endif %}
При append dbt просто добавляет новые строки в конец таблицы. Старые данные не изменяются. Это самая быстрая стратегия - только INSERT, никаких поисков совпадений.
Когда использовать: для данных, которые никогда не изменяются ретроспективно (immutable events). Например, логи запросов к API: каждая строка - уникальное событие, которое не может быть обновлено после записи.
Ограничение: если данные могут поступать с задержкой (late-arriving data), append приведёт к дублям - одно событие будет записано дважды.
insert_overwrite - перезапись партиций
{{ config(
materialized='incremental',
incremental_strategy='insert_overwrite',
partition_by=['event_date'],
file_format='parquet',
) }}
SELECT
event_id,
user_id,
event_date,
channel,
SUM(revenue) AS daily_revenue
FROM {{ ref('stg_events') }}
{% if is_incremental() %}
WHERE event_date >= DATE_ADD(CURRENT_DATE(), -3)
{% endif %}
GROUP BY event_id, user_id, event_date, channel
При insert_overwrite dbt перезаписывает целиком те партиции, которые встречаются в новых данных. Если новые данные содержат event_date = '2024-01-15', dbt удалит всю партицию event_date=2024-01-15 из таблицы и запишет новые строки для этой даты.
Spark выполняет это через INSERT OVERWRITE TABLE ... PARTITION (event_date='2024-01-15') SELECT ....
Важно: в конфиге Spark нужно включить динамическую перезапись партиций:
server_side_parameters:
spark.sql.sources.partitionOverwriteMode: "dynamic"
Когда использовать: для данных с возможными поздними поступлениями (late data), где нужна идемпотентность. Типичный паттерн - пересчёт последних N дней при каждом запуске, что автоматически корректирует ранее загруженные данные.
merge - обновление существующих строк (Delta/Iceberg)
{{ config(
materialized='incremental',
incremental_strategy='merge',
unique_key=['user_id', 'event_date'],
file_format='delta',
partition_by=['event_date'],
) }}
SELECT
user_id,
event_date,
COUNT(DISTINCT session_id) AS sessions,
SUM(revenue) AS total_revenue,
COUNT(*) AS event_count,
CURRENT_TIMESTAMP() AS updated_at
FROM {{ ref('stg_events') }}
{% if is_incremental() %}
WHERE event_date >= DATE_ADD(CURRENT_DATE(), -7)
{% endif %}
GROUP BY user_id, event_date
При merge dbt генерирует SQL MERGE INTO (UPSERT): если строка с таким unique_key уже существует - обновляет её, если не существует - вставляет новую. Это самая мощная стратегия, но требует поддержки MERGE в движке (Delta Lake, Apache Iceberg).
Сгенерированный dbt Spark SQL для merge:
MERGE INTO dbt_prod.fct_user_activity AS DBT_INTERNAL_DEST
USING (
SELECT
user_id,
event_date,
sessions,
total_revenue,
event_count,
updated_at
FROM dbt_prod.fct_user_activity__dbt_tmp
) AS DBT_INTERNAL_SOURCE
ON (
DBT_INTERNAL_SOURCE.user_id = DBT_INTERNAL_DEST.user_id
AND DBT_INTERNAL_SOURCE.event_date = DBT_INTERNAL_DEST.event_date
)
WHEN MATCHED THEN UPDATE SET
sessions = DBT_INTERNAL_SOURCE.sessions,
total_revenue = DBT_INTERNAL_SOURCE.total_revenue,
event_count = DBT_INTERNAL_SOURCE.event_count,
updated_at = DBT_INTERNAL_SOURCE.updated_at
WHEN NOT MATCHED THEN INSERT (
user_id, event_date, sessions, total_revenue, event_count, updated_at
) VALUES (
DBT_INTERNAL_SOURCE.user_id,
DBT_INTERNAL_SOURCE.event_date,
DBT_INTERNAL_SOURCE.sessions,
DBT_INTERNAL_SOURCE.total_revenue,
DBT_INTERNAL_SOURCE.event_count,
DBT_INTERNAL_SOURCE.updated_at
)
Это именно то, что нужно для агрегатов, которые пересчитываются при поступлении поздних данных: строка user_id=100, event_date=2024-01-15 обновится с правильными итоговыми значениями.
Таблица сравнения incremental-стратегий¶
| Стратегия | Поддержка форматов | Дубли? | UPSERT? | Когда использовать |
|---|---|---|---|---|
append |
Parquet, Delta, Iceberg | Возможны | Нет | Immutable события, нет late data |
insert_overwrite |
Parquet, Delta, Iceberg | Нет (партиции) | Нет | Агрегаты по дням, late data |
merge |
Delta, Iceberg | Нет | Да | Агрегаты с обновлением строк |
Полный refresh: когда нужно пересчитать всё¶
Иногда нужно полностью пересчитать incremental-таблицу: изменилась бизнес-логика модели, обнаружили баг, нужно пересчитать исторические данные:
dbt run --select fct_user_activity --full-refresh
С флагом --full-refresh dbt игнорирует блок {% if is_incremental() %} и пересчитывает таблицу полностью, как если бы она была создана впервые. Без этого флага dbt работает в инкрементальном режиме.
Тестирование моделей¶
dbt имеет встроенную систему тестирования, которая запускается командой dbt test. Тесты описываются в YAML-файлах или как SQL-запросы.
Generic тесты (schema tests)¶
# models/staging/schema.yml
version: 2
models:
- name: stg_events
description: "Нормализованные события из Bronze"
columns:
- name: event_id
description: "Уникальный ID события"
tests:
- not_null
- unique
- name: user_id
description: "ID пользователя"
tests:
- not_null
- relationships:
to: ref('stg_users')
field: user_id
- name: event_type
tests:
- not_null
- accepted_values:
values: ['click', 'view', 'purchase', 'signup']
- name: revenue
tests:
- not_null
- name: fct_user_activity
columns:
- name: user_id
tests:
- not_null
- name: event_date
tests:
- not_null
- name: total_revenue
tests:
- not_null
Эти тесты dbt компилирует в SQL-запросы и выполняет. Например, not_null для event_id компилируется в:
SELECT COUNT(*) AS failures
FROM dbt_dev.stg_events
WHERE event_id IS NULL
Если failures > 0 - тест провален.
Custom тесты¶
-- tests/assert_revenue_non_negative.sql
SELECT
user_id,
event_date,
total_revenue
FROM {{ ref('fct_user_activity') }}
WHERE total_revenue < 0
Если этот запрос возвращает строки - тест провален. Это позволяет проверить любую бизнес-логику, которую нельзя выразить через generic тесты.
dbt test --select fct_user_activity
Macros: переиспользуемая логика¶
Макросы в dbt - это Jinja-функции, которые можно использовать в любых SQL-файлах. Они позволяют вынести повторяющиеся паттерны в одно место.
Пример: макрос для стандартизации колонок аудита¶
-- macros/add_audit_columns.sql
{% macro add_audit_columns() %}
CURRENT_TIMESTAMP() AS created_at,
CURRENT_TIMESTAMP() AS updated_at,
'{{ env_var("DBT_RUNNER", "dbt") }}' AS dbt_run_by,
'{{ invocation_id }}' AS dbt_invocation_id
{% endmacro %}
Использование в модели:
SELECT
user_id,
event_date,
SUM(revenue) AS total_revenue,
{{ add_audit_columns() }}
FROM {{ ref('stg_events') }}
GROUP BY user_id, event_date
Пример: макрос для safe-division¶
-- macros/safe_divide.sql
{% macro safe_divide(numerator, denominator, default=0) %}
CASE WHEN {{ denominator }} = 0
THEN {{ default }}
ELSE {{ numerator }} / {{ denominator }}
END
{% endmacro %}
Использование:
SELECT
campaign_id,
clicks,
conversions,
{{ safe_divide('conversions', 'clicks') }} AS conversion_rate
FROM {{ ref('fct_campaign_metrics') }}
Лабораторная практика: Medallion-архитектура на dbt + Spark¶
Задача¶
Построить пайплайн расчёта дневных метрик пользовательской активности. Источник - сырые логи событий в MinIO (Parquet). Цель - витрина fct_user_activity в Delta Lake с инкрементальным обновлением по стратегии merge.
Шаг 1: Инициализация проекта¶
pip install dbt-spark[PyHive] dbt-core
dbt init analytics_lakehouse
cd analytics_lakehouse
Дерево создаётся автоматически. Заполним profiles.yml:
# ~/.dbt/profiles.yml
analytics_lakehouse:
target: dev
outputs:
dev:
type: spark
method: session
schema: dbt_dev
server_side_parameters:
spark.sql.shuffle.partitions: "50"
spark.sql.adaptive.enabled: "true"
spark.sql.extensions: "io.delta.sql.DeltaSparkSessionExtension"
spark.sql.catalog.spark_catalog: "org.apache.spark.sql.delta.catalog.DeltaCatalog"
method: session - режим для локального тестирования без Thrift Server. dbt запускает SparkSession прямо в Python-процессе.
Шаг 2: Описание источников¶
# models/staging/sources.yml
version: 2
sources:
- name: bronze
description: "Bronze слой - сырые данные"
database: bronze
schema: raw
tables:
- name: events
identifier: events
description: "Сырые события пользователей"
loaded_at_field: event_ts
freshness:
warn_after: {count: 2, period: hour}
error_after: {count: 6, period: hour}
columns:
- name: event_id
tests: [not_null, unique]
- name: user_id
tests: [not_null]
- name: event_ts
tests: [not_null]
- name: users
identifier: users
description: "Справочник пользователей из CDC"
loaded_at_field: updated_at
Шаг 3: Staging-модели (Silver)¶
-- models/staging/stg_events.sql
{{ config(
materialized='view',
tags=['staging', 'events']
) }}
SELECT
event_id,
CAST(user_id AS BIGINT) AS user_id,
session_id,
event_type,
CAST(event_ts AS TIMESTAMP) AS event_ts,
CAST(event_ts AS DATE) AS event_date,
COALESCE(CAST(revenue AS DECIMAL(12,4)), 0.0) AS revenue,
LOWER(TRIM(COALESCE(channel, 'unknown'))) AS channel,
LOWER(TRIM(COALESCE(device_type, 'unknown'))) AS device_type,
CURRENT_TIMESTAMP() AS dbt_loaded_at
FROM {{ source('bronze', 'events') }}
WHERE event_id IS NOT NULL
AND user_id IS NOT NULL
AND event_ts IS NOT NULL
-- models/staging/stg_users.sql
{{ config(
materialized='view',
tags=['staging', 'users']
) }}
SELECT
CAST(user_id AS BIGINT) AS user_id,
LOWER(TRIM(email)) AS email,
COALESCE(first_name, 'Unknown') AS first_name,
COALESCE(last_name, 'Unknown') AS last_name,
CAST(created_at AS TIMESTAMP) AS created_at,
CAST(created_at AS DATE) AS signup_date,
CASE
WHEN account_type IN ('premium', 'enterprise') THEN TRUE
ELSE FALSE
END AS is_paid_user,
COALESCE(country_code, 'XX') AS country_code,
CURRENT_TIMESTAMP() AS dbt_loaded_at
FROM {{ source('bronze', 'users') }}
WHERE user_id IS NOT NULL
AND email IS NOT NULL
Шаг 4: Intermediate-модели¶
-- models/intermediate/int_user_sessions.sql
{{ config(materialized='ephemeral') }}
SELECT
user_id,
session_id,
event_date,
MIN(event_ts) AS session_start,
MAX(event_ts) AS session_end,
DATEDIFF(SECOND, MIN(event_ts), MAX(event_ts)) AS session_duration_sec,
COUNT(*) AS events_in_session,
SUM(revenue) AS session_revenue,
COUNT(DISTINCT event_type) AS unique_event_types,
MAX(channel) AS channel,
MAX(device_type) AS device_type
FROM {{ ref('stg_events') }}
GROUP BY user_id, session_id, event_date
Шаг 5: Gold-витрина с incremental merge¶
-- models/marts/fct_user_activity.sql
{{ config(
materialized='incremental',
incremental_strategy='merge',
unique_key=['user_id', 'event_date'],
file_format='delta',
partition_by=['event_date'],
tags=['gold', 'daily', 'user-metrics']
) }}
WITH sessions AS (
SELECT * FROM {{ ref('int_user_sessions') }}
{% if is_incremental() %}
WHERE event_date >= DATE_ADD(CURRENT_DATE(), -3)
{% endif %}
),
events AS (
SELECT * FROM {{ ref('stg_events') }}
{% if is_incremental() %}
WHERE event_date >= DATE_ADD(CURRENT_DATE(), -3)
{% endif %}
),
user_info AS (
SELECT user_id, email, is_paid_user, country_code, signup_date
FROM {{ ref('stg_users') }}
)
SELECT
e.user_id,
u.email,
u.is_paid_user,
u.country_code,
u.signup_date,
e.event_date,
COUNT(DISTINCT e.session_id) AS sessions,
SUM(s.session_duration_sec) AS total_session_sec,
AVG(s.session_duration_sec) AS avg_session_sec,
SUM(e.revenue) AS total_revenue,
COUNT(*) AS event_count,
COUNT(DISTINCT e.event_type) AS unique_event_types,
COUNT(DISTINCT s.channel) AS channels_used,
MAX(s.channel) AS primary_channel,
CURRENT_TIMESTAMP() AS updated_at
FROM events e
LEFT JOIN user_info u ON e.user_id = u.user_id
LEFT JOIN sessions s ON e.user_id = s.user_id
AND e.session_id = s.session_id
AND e.event_date = s.event_date
GROUP BY
e.user_id, u.email, u.is_paid_user, u.country_code, u.signup_date, e.event_date
Обратите внимание: блок {% if is_incremental() %} применён дважды - к CTE sessions и к CTE events. Это ключевой момент: при инкрементальном запуске нужно ограничить не только финальный SELECT, но и промежуточные шаги, иначе int_user_sessions (ephemeral) будет вычислен по всему датасету.
Шаг 6: YAML-схема для marts¶
# models/marts/schema.yml
version: 2
models:
- name: fct_user_activity
description: |
Дневные метрики активности пользователей.
Обновляется инкрементально через MERGE по (user_id, event_date).
Окно обновления: последние 3 дня (для компенсации late data).
columns:
- name: user_id
description: "ID пользователя"
tests: [not_null]
- name: event_date
description: "Дата события"
tests: [not_null]
- name: total_revenue
description: "Суммарная выручка за день"
tests: [not_null]
- name: sessions
tests: [not_null]
Шаг 7: Запуск пайплайна¶
dbt run --select staging.* marts.fct_user_activity
dbt run --select fct_user_activity+
dbt run --select +fct_user_activity+
dbt run --select fct_user_activity --full-refresh
dbt test --select fct_user_activity
Синтаксис выборки моделей:
model_name- только эта модельmodel_name+- эта модель и все downstream-зависимости+model_name- эта модель и все upstream-зависимости+model_name+- всё дерево вокруг моделиtag:gold- все модели с тегомgold
Анализ сгенерированного SQL в логах¶
После запуска dbt сохраняет скомпилированные запросы и подробные логи:
cat target/compiled/analytics_lakehouse/models/marts/fct_user_activity.sql
cat logs/dbt.log | grep "MERGE INTO"
В логах можно найти полный MERGE INTO запрос, который dbt отправил в Spark. Это позволяет:
- Убедиться, что фильтрация по датам работает корректно
- Скопировать и выполнить запрос вручную для отладки
- Проанализировать план выполнения через
EXPLAINв Spark SQL
dbt compile --select fct_user_activity
cat target/compiled/analytics_lakehouse/models/marts/fct_user_activity.sql | \
spark-sql --master spark://localhost:7077 -e "EXPLAIN FORMAT=FORMATTED $(cat -)"
Шаг 8: Документация¶
dbt docs generate
dbt docs serve --port 8080
Откройте http://localhost:8080 в браузере. Вы увидите:
- Каталог всех моделей с описаниями
- Lineage Graph для
fct_user_activity- интерактивная схема зависимостей - Статус тестов для каждой колонки
- Скомпилированный SQL для каждой модели
Переменные и окружения¶
dev vs prod изоляция¶
Каждый разработчик работает в своей "схеме" (базе данных). По умолчанию dbt использует schema из profiles.yml - например, dbt_alice, dbt_bob. Это обеспечивает изоляцию: каждый видит свои таблицы, не мешая другим.
Кастомизация имён схем через макрос:
-- macros/generate_schema_name.sql
{% macro generate_schema_name(custom_schema_name, node) -%}
{%- set default_schema = target.schema -%}
{%- if custom_schema_name is none -%}
{{ default_schema }}
{%- else -%}
{{ default_schema }}_{{ custom_schema_name | trim }}
{%- endif -%}
{%- endmacro %}
С этим макросом модель с +schema: staging в dev создаётся как dbt_alice_staging, в prod - как dbt_prod_staging.
CI/CD: state-aware выполнение¶
Мощная функция dbt - запуск только изменившихся моделей:
git checkout main
dbt compile
cp -r target/ target_prod/
git checkout feature/new-metric
dbt run --select state:modified+ \
--state target_prod/
dbt сравнивает скомпилированный SQL текущей ветки с эталонным (target_prod) и запускает только модели, которые изменились, плюс все их downstream-зависимости. Это драматически сокращает время CI: вместо пересчёта 50 моделей запускаются только 3 изменившиеся.
Антипаттерны¶
Антипаттерн 1: Гигантская монолитная модель¶
-- models/marts/everything.sql ← ПЛОХО
SELECT
e.event_id,
e.user_id,
u.email,
u.country,
c.campaign_name,
c.budget,
p.product_name,
p.category,
s.session_duration,
...
(500 строк SQL)
FROM bronze.raw_events e
LEFT JOIN bronze.raw_users u ON e.user_id = u.user_id
LEFT JOIN bronze.raw_campaigns c ON e.campaign_id = c.campaign_id
LEFT JOIN bronze.raw_products p ON e.product_id = p.product_id
LEFT JOIN ...
Проблемы: нет повторного использования, нет lineage, нет возможности протестировать части независимо, невозможно кэшировать промежуточные результаты. Правило: одна модель - одна бизнес-сущность или один чёткий шаг трансформации.
Антипаттерн 2: Hardcode имён таблиц¶
-- ПЛОХО: хардкод путей
SELECT * FROM prod_db.gold.fct_user_activity
JOIN bronze.raw.events ON ...
Вместо ref() и source() - ручные имена. Ломается при смене окружения, не отображается в lineage.
Антипаттерн 3: Сверхизбыточные views для тяжёлых агрегаций¶
-- models/staging/stg_daily_revenue.sql
{{ config(materialized='view') }}
SELECT
event_date,
SUM(revenue) AS total_revenue
FROM {{ source('bronze', 'events') }}
GROUP BY event_date
View с тяжёлой агрегацией над миллиардами строк - каждый запрос к этой view будет сканировать весь Bronze-слой. Для тяжёлых агрегаций нужна материализация table или incremental.
Антипаттерн 4: Missing incremental filter¶
-- models/marts/fct_revenue.sql
{{ config(materialized='incremental') }}
SELECT * FROM {{ ref('stg_events') }}
{% if is_incremental() %}
-- Забыли добавить WHERE условие!
{% endif %}
Если блок is_incremental() пустой, при инкрементальном запуске обрабатываются все данные - смысл инкрементальности теряется. Всегда добавляйте WHERE-условие в блок is_incremental().
Антипаттерн 5: unique_key из одного поля для агрегатов¶
{{ config(
materialized='incremental',
unique_key='user_id',
incremental_strategy='merge'
) }}
SELECT
user_id,
event_date,
SUM(revenue) AS total_revenue
FROM {{ ref('stg_events') }}
GROUP BY user_id, event_date
unique_key='user_id' - только по user_id. Это значит, для каждого пользователя в таблице будет только одна строка. Предыдущий день будет перезаписан данными нового дня! Правильно: unique_key=['user_id', 'event_date'].
Антипаттерн 6: Игнорировать Spark-специфичный тюнинг¶
Модель dbt выполняется как Spark SQL-запрос. Все принципы оптимизации - shuffle partitions, broadcast join, partition pruning - применимы и здесь. Добавляйте hints прямо в SQL:
SELECT /*+ BROADCAST(u) */
e.user_id,
u.email,
SUM(e.revenue) AS total_revenue
FROM {{ ref('stg_events') }} e
JOIN {{ ref('stg_users') }} u ON e.user_id = u.user_id
GROUP BY e.user_id, u.email
Или используйте server_side_parameters в profiles.yml для настройки shuffle.partitions.
Чеклист¶
Структура проекта:
- Staging-модели: только
view, один источник - одна модель, переименование и кастинг типов - Intermediate-модели:
ephemeralдля промежуточной логики,tableдля повторно используемых - Mart-модели:
incrementalдля больших данных,tableдля малых Gold-витрин - Каждая модель имеет описание в schema.yml
- Все ключевые колонки покрыты тестами (not_null, unique, relationships)
Инкрементальные модели:
- Всегда есть фильтр в блоке
{% if is_incremental() %} unique_keyвключает все поля, определяющие уникальность строки- Окно пересчёта перекрывает N дней для компенсации late data
- Выбрана правильная стратегия: append/insert_overwrite/merge
- Для
mergeиспользуется Delta Lake или Iceberg
Источники:
- Все внешние таблицы зарегистрированы в sources.yml
- Задана freshness для критических источников
- Используется
{{ source() }}вместо хардкода
Домашнее задание¶
Дана витрина продаж, пересчитываемая полностью при каждом запуске:
-- models/marts/fct_sales_daily.sql
{{ config(materialized='table') }}
SELECT
order_id,
user_id,
product_id,
order_date,
CAST(quantity AS INT) AS quantity,
CAST(unit_price AS DECIMAL(12,4)) AS unit_price,
quantity * unit_price AS line_revenue,
status,
updated_at
FROM {{ ref('stg_orders') }}
WHERE status IN ('completed', 'shipped')
Пайплайн работает 2 часа на терабайтном датасете заказов. Размер прироста - один день ≈ 500 МБ.
Задание 1. Переведите модель в incremental с insert_overwrite по order_date. Обоснуйте выбор именно этой стратегии (а не append или merge).
Задание 2. Напишите правильный фильтр в блоке {% if is_incremental() %} с учётом того, что заказы могут обновляться (поле status может измениться) в течение 7 дней после создания.
Задание 3. Добавьте в конфиг правильные настройки для Spark:
- Формат файлов (Parquet или Delta - аргументируйте выбор)
- Партиционирование
- Spark-параметры в profiles.yml для оптимальной работы (shuffle.partitions для 500 МБ/день при кластере 16 × 4 ядра)
Задание 4. Откройте Spark UI после первого полного запуска и после первого инкрементального запуска. Что изменилось в метриках? Сколько данных сканируется при инкрементальном запуске по сравнению с полным? Сформулируйте ответ с конкретными числами из Spark UI (Shuffle Read Size, Input Size, Task Duration).