dbt Core: models, refs, sources и materialization types (table, view, incremental)

dbt Core превращает SELECT-запросы в управляемый DAG трансформаций поверх Spark SQL. Разбираем sources, ref(), стратегии материализации и incremental-модели для lakehouse.

platform

Зачем 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, он:

  1. Скомпилирует Jinja-шаблоны, подставив реальные имена таблиц
  2. Создаст VIEW (потому что materialized='view') в базе данных: CREATE OR REPLACE VIEW dbt_dev.stg_users AS SELECT ...
  3. Зарегистрирует зависимость: 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. На практике это означает:

  1. dbt создаёт временную таблицу dbt_dev.fct_revenue_summary__dbt_tmp
  2. Выполняет SELECT в эту временную таблицу
  3. Переименовывает временную таблицу в 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).