dbt Tests и Source Freshness: generic tests, singular tests, data contracts

dbt Tests и Source Freshness: generic tests, singular tests, data contracts

platform

Концепция Data Quality в Data Lakehouse

Проблема «мусор на входе - мусор на выходе»

В классических СУБД (PostgreSQL, MySQL) база данных физически не позволяет записать некорректные данные: ограничение NOT NULL отклоняет пустые значения, PRIMARY KEY блокирует дубликаты, FOREIGN KEY запрещает несуществующие ссылки. Система гарантирует целостность на уровне движка - это называется принудительная схема (schema-on-write).

В мире Data Lakehouse всё устроено иначе. Apache Parquet - это файловый формат без встроенного сервера. Никто не проверяет, что вы записываете в файл. Spark читает данные с диска и интерпретирует их согласно схеме - это schema-on-read. Если Bronze-слой содержит дубликаты из CDC-конвейера, нулевые значения там, где должны быть суммы, или отрицательные цены - Spark прочитает это без каких-либо предупреждений. Данные просто пройдут через все трансформации и окажутся в Gold-витрине, которую видят аналитики и топ-менеджмент.

Проблема «мусор на входе - мусор на выходе» (GIGO, Garbage In - Garbage Out) в Data Engineering проявляется коварно: пайплайн работает, метрики считаются, дашборды обновляются. Только цифры неправильные. Это хуже, чем сломанный пайплайн: сломанный пайплайн заметят сразу, а тихая порча данных может жить месяцами.

Конкретный пример: CDC-репликация из PostgreSQL в MinIO работает с задержкой. Вместо свежих данных витрина fct_orders содержит данные двухдневной давности, но никто об этом не знает. Финансовый аналитик строит отчёт по продажам за вчера, получает заниженные цифры и начинает выяснять, «куда делись заказы». Несколько часов работы нескольких человек - чтобы выяснить, что данные просто устарели.

Shift-Left: тестирование на этапе трансформации

Концепция Shift-Left в Data Engineering означает перенос контроля качества как можно левее по пайплайну - то есть как можно раньше. Лучше поймать проблему на Bronze-слое, чем узнать о ней из жалоб бизнеса.

                  Стоимость исправления растёт →→→

[Ingestion]  →  [Bronze]  →  [Silver]  →  [Gold]  →  [Report]
    ↑              ↑             ↑            ↑           ↑
  Дёшево         Дёшево       Дороже      Дорого    Очень дорого

dbt позволяет встроить проверки прямо в пайплайн трансформации. Если тест не прошёл - выполнение можно остановить до того, как некачественные данные попадут вниз по слоям. dbt build (в отличие от dbt run) делает именно это: запускает трансформации и тесты в единой DAG, прерывая выполнение при провале любого теста.

Специфика Spark: почему тестирование обязательно

Несколько особенностей Spark-экосистемы делают контроль качества особенно важным:

Schema-on-read и Parquet. Когда вы пишете spark.read.parquet("s3a://bucket/data/"), Spark берёт схему из метаданных файла. Если разные файлы имеют разные схемы (такое бывает при схема-эволюции), Spark сделает слияние схем - но без гарантий относительно NULL-значений там, где раньше были данные.

Distributed ingestion. Данные могут приходить из десятков источников параллельно. Один источник прислал дубликаты, другой - пустые строки, третий - значения в неправильной валюте. Spark загрузит всё без жалоб.

CDC и идемпотентность. При повторном запуске CDC-конвейера (сеть упала, джоб перезапустился) часть событий может задублироваться. Без проверки unique это незаметно.

Отсутствие физических ограничений. Delta Lake и Iceberg начали поддерживать ограничения (NOT NULL, CHECK) - но это относительно новая функциональность, доступная не везде. dbt-тесты работают независимо от формата и движка.

Архитектура тестирования в dbt

dbt реализует тестирование через SQL. Каждый тест - это SQL-запрос, который возвращает строки-нарушители. Если запрос вернул 0 строк - тест прошёл. Если вернул хотя бы одну строку - тест упал.

-- Пример: тест unique для колонки order_id
SELECT order_id
FROM my_schema.fct_orders
WHERE order_id IS NOT NULL
GROUP BY order_id
HAVING COUNT(*) > 1

-- Если этот SELECT вернул строки - есть дубликаты → тест упал
-- Если SELECT вернул 0 строк - дублей нет → тест прошёл

Это принципиально отличается от assertion-подхода в unit-тестировании: нет assertEqual, нет assertTrue. Только SQL и количество строк в результате.

Такой подход имеет важное свойство: тесты запускаются на Spark, а значит масштабируются на терабайты так же, как основные трансформации. Нет нужды сэмплировать данные - проверка идёт по всему объёму.


Source Freshness: контроль актуальности данных

Что такое свежесть данных и почему это важно

Source Freshness - это автоматическая проверка того, насколько свежими являются данные в источниках (Bronze-слой). По сути, это мониторинг pipeline SLA: если данные не обновлялись дольше допустимого - что-то сломалось.

Типичные сценарии, когда freshness-проверка спасает:

  • ETL завис: Airflow-джоб завис в состоянии Running, но данные не льются. Spark-job завершился с ошибкой, но оркестратор не заметил.
  • Kafka lag: Топик накопил лаг, consumer не успевает, Bronze обновляется с задержкой.
  • Сетевой сбой: Репликация из PostgreSQL остановилась из-за сетевой проблемы.
  • Изменение расписания: Команда бэкенда изменила частоту выгрузки без уведомления.

Без freshness-мониторинга все эти проблемы обнаруживаются только когда аналитики жалуются на устаревшие данные - то есть слишком поздно.

Конфигурация sources.yml

Freshness настраивается в файле sources.yml (или внутри schema.yml) - там же, где описываются источники данных:

# models/staging/sources.yml
version: 2

sources:
  - name: bronze_transactions
    description: >
      Bronze-слой с транзакциями из PostgreSQL через CDC-репликацию.
      Обновляется каждые 15 минут через Debezium + Kafka.
    database: hive_metastore
    schema: bronze

    # Колонка, по которой измеряем свежесть.
    # Это должен быть timestamp момента загрузки в Bronze,
    # а не бизнес-timestamp (created_at транзакции).
    loaded_at_field: _ingested_at

    # SLA по умолчанию для всех таблиц источника
    freshness:
      warn_after: {count: 1, period: hour}
      error_after: {count: 3, period: hours}

    tables:
      - name: raw_orders
        description: Сырые заказы из PostgreSQL orders table
        # Переопределяем freshness для этой конкретной таблицы
        freshness:
          warn_after: {count: 30, period: minutes}
          error_after: {count: 2, period: hours}

      - name: raw_order_items
        description: Позиции заказов - обновляется вместе с raw_orders
        # Наследует freshness от источника: warn=1h, error=3h

      - name: raw_users
        description: Профили пользователей - медленно меняются
        # Для справочников допустима большая задержка
        freshness:
          warn_after: {count: 6, period: hours}
          error_after: {count: 24, period: hours}

      - name: raw_product_catalog
        description: Каталог товаров
        # Полностью отключаем проверку freshness
        freshness: null

Параметры warn_after и error_after задают SLA:

  • warn_after - предупреждение (тест завершается с WARN, пайплайн не прерывается)
  • error_after - ошибка (тест завершается с ERROR, пайплайн прерывается при dbt build)
  • period - единица времени: minute, hour, day

Как dbt проверяет freshness

При запуске dbt source freshness dbt компилирует для каждой таблицы быстрый агрегирующий запрос:

-- Что dbt генерирует для raw_orders:
SELECT
    MAX(_ingested_at) AS max_loaded_at,
    CURRENT_TIMESTAMP() AS snapshotted_at
FROM hive_metastore.bronze.raw_orders

Это важная деталь: не полный скан таблицы, а только SELECT MAX(timestamp). На партиционированной таблице Spark выполнит этот запрос за секунды, сканируя только метаданные партиций (если _ingested_at - колонка партиционирования). На S3 это может быть буквально один LIST-запрос.

dbt сравнивает результат с текущим временем:

Текущее время: 2024-01-15 14:30:00
MAX(_ingested_at): 2024-01-15 12:45:00
Разница: 1 час 45 минут

warn_after: 30 минут  → WARN (превышено)
error_after: 2 часа   → OK (не превышено)
→ Статус: WARN

Запуск и интерпретация результатов

# Запустить freshness-проверки для всех источников
dbt source freshness

# Только для конкретного источника
dbt source freshness --select source:bronze_transactions

# Вывод в файл для мониторинга
dbt source freshness --output sources.json

Пример вывода:

Running with dbt=1.7.0
Found 3 sources, 0 freshness errors

14:30:01  Freshness of bronze_transactions.raw_orders ........ [WARN in 0.84s]
14:30:02  Freshness of bronze_transactions.raw_order_items ... [pass in 0.71s]
14:30:03  Freshness of bronze_transactions.raw_users ......... [pass in 0.62s]

Done. PASS=2 WARN=1 ERROR=0 SKIP=0 TOTAL=3

Результат записывается в target/sources.json - этот файл можно подать в системы мониторинга (Grafana, Datadog) для построения алертов.

Freshness для партиционированных таблиц

Если Bronze-таблица партиционирована по дате, freshness-запрос автоматически выигрывает от pruning:

tables:
  - name: raw_events
    loaded_at_field: event_date  # Поле партиционирования
    freshness:
      warn_after: {count: 6, period: hours}
-- dbt генерирует:
SELECT MAX(event_date) AS max_loaded_at, CURRENT_TIMESTAMP() AS snapshotted_at
FROM bronze.raw_events

Spark выполнит MAX(event_date) через метаданные Metastore - список партиций - без чтения файлов. Буквально мгновенно.

Freshness и orchestration

Freshness-проверку можно встроить в Airflow-пайплайн как gate перед запуском трансформаций:

# Airflow DAG
from airflow.operators.bash import BashOperator

check_freshness = BashOperator(
    task_id='check_source_freshness',
    bash_command='cd /opt/dbt && dbt source freshness --select source:bronze_transactions',
)

run_transformations = BashOperator(
    task_id='run_dbt_transformations',
    bash_command='cd /opt/dbt && dbt run --select marts.*',
)

# Трансформации запустятся только если freshness прошёл
check_freshness >> run_transformations

При ERROR в freshness Airflow-таск завершится с ненулевым кодом, следующие таски не запустятся, и команда получит алерт. Это предотвращает загрузку устаревших данных в Gold-слой.


Generic Tests: декларативная валидация

Что такое generic tests

Generic tests (или schema tests) - это встроенные в dbt переиспользуемые проверки, которые конфигурируются декларативно в YAML. Вам не нужно писать SQL-запрос вручную - достаточно объявить, что order_id должен быть уникальным, и dbt сам сгенерирует нужный SQL.

dbt Core включает четыре стандартных generic теста:

  • not_null - колонка не содержит NULL
  • unique - все значения уникальны
  • accepted_values - значения из заданного допустимого набора
  • relationships - ссылочная целостность между таблицами

Конфигурация в schema.yml

# models/marts/schema.yml
version: 2

models:
  - name: fct_orders
    description: >
      Финансовая витрина заказов. Одна строка - один заказ.
      Обновляется инкрементально каждый час.

    config:
      # Строгий режим: тест-ошибка блокирует downstream модели
      severity: error

    columns:
      - name: order_id
        description: Уникальный идентификатор заказа из PostgreSQL
        tests:
          - not_null
          - unique

      - name: user_id
        description: Ссылка на пользователя
        tests:
          - not_null
          - relationships:
              to: ref('dim_users')
              field: user_id

      - name: status
        description: Статус заказа
        tests:
          - not_null
          - accepted_values:
              values: ['pending', 'confirmed', 'shipped', 'delivered', 'cancelled', 'refunded']

      - name: total_amount
        description: Итоговая сумма заказа в рублях
        tests:
          - not_null

      - name: event_date
        description: Дата создания заказа (partition key)
        tests:
          - not_null

Тест not_null

Самый простой и самый важный тест. Проверяет, что в колонке нет значений NULL.

Что генерирует dbt:

-- target/compiled/.../not_null_fct_orders_order_id.sql
SELECT order_id
FROM my_schema.fct_orders
WHERE order_id IS NULL

Если строки с NULL существуют - тест падает. Применяйте not_null ко всем бизнес-ключам и обязательным полям: ID-шники, суммы, даты событий. Это минимальный baseline, без которого нельзя доверять данным.

Частая ошибка - не тестировать поля, которые «не могут быть NULL по логике». Могут. CDC-репликация прислала событие с пустым user_id - Spark записал NULL. Без теста это обнаружат аналитики, а не инженеры.

Тест unique

Проверяет, что все значения в колонке уникальны (NULL-значения игнорируются).

-- Генерируемый SQL:
SELECT order_id
FROM my_schema.fct_orders
WHERE order_id IS NOT NULL
GROUP BY order_id
HAVING COUNT(*) > 1

unique - это проверка бизнес-инварианта: «одна строка = один объект». Если fct_orders содержит дубликаты по order_id, агрегация по заказам даст задвоенные суммы. Это классический источник расхождений между аналитическими отчётами.

Применяйте unique ко всем суррогатным ключам и бизнес-идентификаторам. Если таблица представляет snapshot (история изменений одного объекта), то unique применяется к составному ключу.

Тест accepted_values

Проверяет, что все значения колонки принадлежат заданному множеству. Используется для enum-полей: статусы, категории, коды стран.

- name: status
  tests:
    - accepted_values:
        values: ['pending', 'confirmed', 'shipped', 'delivered', 'cancelled', 'refunded']
        # По умолчанию quote=true (значения в кавычках - строки)
        # Для числовых значений:
        # quote: false
        # values: [1, 2, 3, 4]
-- Генерируемый SQL:
SELECT DISTINCT status
FROM my_schema.fct_orders
WHERE status NOT IN ('pending', 'confirmed', 'shipped', 'delivered', 'cancelled', 'refunded')
  AND status IS NOT NULL

Особенно важен при интеграции с внешними системами: если бэкенд добавил новый статус ('returned'), но аналитики его не знают - тест поймает это немедленно. Новый статус не «потеряется» в подсчётах по статусам, а вызовет ошибку в пайплайне.

Тест relationships

Проверяет ссылочную целостность: каждое значение в колонке A таблицы X должно существовать в колонке B таблицы Y. Аналог FOREIGN KEY, но проверяется runtime, а не схемой.

- name: user_id
  tests:
    - relationships:
        to: ref('dim_users')
        field: user_id
-- Генерируемый SQL:
SELECT order_id
FROM my_schema.fct_orders
WHERE user_id IS NOT NULL
  AND user_id NOT IN (
      SELECT user_id FROM my_schema.dim_users WHERE user_id IS NOT NULL
  )

Обнаруживает «висячие ссылки» - заказы от пользователей, которых нет в dim_users. Это типичная проблема при асинхронной репликации: заказ пришёл раньше, чем профиль пользователя.

Производительность: NOT IN (subquery) на больших таблицах может быть медленным. Spark выполнит join двух таблиц. Для таблиц объёмом > 100 GB используйте severity: warn или выборочную проверку через where (об этом далее).

Параметры тестов

Severity: warn vs error

По умолчанию все тесты - error (проваленный тест прерывает dbt build). Можно снизить до warn:

- name: user_id
  tests:
    - relationships:
        to: ref('dim_users')
        field: user_id
        severity: warn   # Предупреждение, не ошибка

Используйте warn для проверок, которые нарушаются по объективным причинам (например, relationships при асинхронной загрузке) или для новых тестов в период «обкатки».

where: фильтрация данных

Тест можно применять не ко всем строкам, а к подмножеству:

- name: total_amount
  tests:
    - not_null:
        where: "status != 'cancelled'"
        # NULL в total_amount допустим только для отменённых заказов

where - это мощный инструмент для избежания ложных срабатываний без отключения теста полностью. Также используется для оптимизации: тестировать только свежие данные.

- name: order_id
  tests:
    - unique:
        where: "event_date >= current_date() - interval 7 days"
        # Уникальность проверяем только за последние 7 дней
        # Для исторических данных уникальность уже была проверена

store_failures: сохранение проблемных строк

При падении теста нужно знать, какие конкретно строки нарушают инвариант. Параметр store_failures создаёт отдельную таблицу с проблемными записями:

models:
  - name: fct_orders
    columns:
      - name: order_id
        tests:
          - unique:
              store_failures: true
              # Создаст таблицу: my_schema_dbt_test__audit.unique_fct_orders_order_id

После падения теста:

-- Смотрим дубликаты напрямую:
SELECT * FROM my_schema_dbt_test__audit.unique_fct_orders_order_id
LIMIT 100;

Это гораздо удобнее, чем пытаться воспроизвести проблему вручную.

Composite uniqueness: составные бизнес-ключи

Когда уникальность определяется не одной колонкой, а комбинацией:

models:
  - name: fct_order_items
    description: Позиции заказов. Уникальность - комбинация order_id + item_id.
    tests:
      # Тест на уровне модели (не колонки)
      - unique:
          column_name: "CONCAT(order_id, '|', item_id)"
          # Или через dbt-utils (подробнее ниже)

Лучший вариант - использовать dbt_utils.unique_combination_of_columns:

    tests:
      - dbt_utils.unique_combination_of_columns:
          combination_of_columns:
            - order_id
            - item_id

Пакет dbt-utils: расширенные generic tests

dbt-utils - официальный пакет с дополнительными generic tests. Установка:

# packages.yml
packages:
  - package: dbt-labs/dbt_utils
    version: 1.1.1
dbt deps  # Установить пакеты

Наиболее полезные тесты из dbt-utils:

models:
  - name: fct_orders
    tests:
      # Уникальность составного ключа
      - dbt_utils.unique_combination_of_columns:
          combination_of_columns: [order_id, item_id]

      # Колонка возрастает монотонно (для суррогатных ключей)
      - dbt_utils.monotonically_increasing_id:
          column_name: surrogate_key

    columns:
      - name: total_amount
        tests:
          # Произвольное условие SQL (самый гибкий generic test)
          - dbt_utils.expression_is_true:
              expression: ">= 0"
              # Проверяет: total_amount >= 0

          # Значение в диапазоне
          - dbt_utils.accepted_range:
              min_value: 0
              max_value: 10000000  # Максимальный заказ 10 млн руб

      - name: created_at
        tests:
          # Дата не в будущем
          - dbt_utils.expression_is_true:
              expression: "<= current_timestamp()"

expression_is_true - самый универсальный тест в dbt-utils. Он проверяет, что для каждой строки SQL-выражение истинно. Если нет - строка попадает в результат (нарушитель).


Singular Tests: кастомная бизнес-логика

Почему generic tests недостаточно

Generic tests отлично справляются со структурными инвариантами: не-NULL, уникальность, допустимые значения. Но бизнес-логика часто требует проверок, которые нельзя выразить одним условием на одной колонке:

  • «Сумма скидки не может превышать стоимость позиции»
  • «Возврат не может быть создан позже даты гарантии»
  • «Сумма транзакций по пользователю должна совпадать с изменением его баланса»
  • «Событие завершения заказа не может предшествовать событию его создания»

Для таких проверок существуют singular tests - кастомные SQL-файлы в директории tests/.

Структура singular теста

Singular test - это обычный .sql файл. Правило одно: если SELECT возвращает строки - тест упал.

-- tests/assert_discount_not_exceeds_total.sql
-- Тест: скидка не может превышать полную стоимость позиции

SELECT
    order_id,
    item_id,
    total_amount,
    discount_amount,
    (discount_amount - total_amount) AS overcharge
FROM {{ ref('fct_order_items') }}
WHERE discount_amount > total_amount

Логика: ищем строки-нарушители. Если есть хотя бы одна позиция, где скидка больше суммы - тест падает. dbt покажет количество нарушителей в логе.

Тест запускается так же, как generic:

dbt test --select test_type:singular
dbt test --select assert_discount_not_exceeds_total

Singular тест с параметрами

Singular тесты поддерживают Jinja, как и обычные модели:

-- tests/assert_no_future_orders.sql
-- Заказы не могут быть созданы в будущем (с допуском в 1 час на clock skew)

{% set tolerance_hours = var('clock_skew_tolerance_hours', 1) %}

SELECT
    order_id,
    created_at,
    CURRENT_TIMESTAMP() AS check_time,
    TIMESTAMPDIFF(HOUR, CURRENT_TIMESTAMP(), created_at) AS hours_in_future
FROM {{ ref('fct_orders') }}
WHERE created_at > CURRENT_TIMESTAMP() + INTERVAL {{ tolerance_hours }} HOURS
dbt test --vars '{"clock_skew_tolerance_hours": 2}'

Бизнес-правило: временна́я последовательность событий

-- tests/assert_event_ordering.sql
-- Событие delivered не может предшествовать событию confirmed

SELECT
    order_id,
    confirmed_at,
    delivered_at,
    TIMESTAMPDIFF(MINUTE, confirmed_at, delivered_at) AS minutes_diff
FROM {{ ref('fct_orders') }}
WHERE status = 'delivered'
  AND delivered_at < confirmed_at

Этот тест обнаружит аномалии в CDC-потоке: если события пришли не по порядку (частая проблема в Kafka с несколькими партициями), временна́я метка в витрине может оказаться перевёрнутой.

Бизнес-правило: сверка балансов

Один из наиболее ценных типов тестов - reconciliation test (сверка контрольных сумм):

-- tests/assert_transaction_balance_reconciliation.sql
-- Сверка: изменение баланса пользователя должно совпадать
-- с суммой транзакций за тот же период

WITH balance_changes AS (
    SELECT
        user_id,
        SUM(amount) AS total_transactions
    FROM {{ ref('fct_transactions') }}
    WHERE event_date >= CURRENT_DATE() - INTERVAL 1 DAY
    GROUP BY user_id
),

balance_movements AS (
    SELECT
        user_id,
        SUM(CASE WHEN event_type = 'credit' THEN amount ELSE -amount END) AS net_movement
    FROM {{ ref('fct_balance_movements') }}
    WHERE event_date >= CURRENT_DATE() - INTERVAL 1 DAY
    GROUP BY user_id
),

reconciliation AS (
    SELECT
        t.user_id,
        t.total_transactions,
        b.net_movement,
        ABS(t.total_transactions - b.net_movement) AS discrepancy
    FROM balance_changes t
    JOIN balance_movements b ON t.user_id = b.user_id
)

SELECT *
FROM reconciliation
WHERE discrepancy > 0.01  -- Допуск 1 копейка на ошибки округления

Такой тест - финансовая страховка. Он гарантирует, что две независимые таблицы согласуются между собой. Если где-то произошло задвоение транзакции или потеря события - тест немедленно обнаружит расхождение.

Singular тест для инкрементальных моделей

При инкрементальном обновлении важно проверить, что новые данные правильно приклеились к историческим:

-- tests/assert_no_gaps_in_daily_data.sql
-- Проверка: нет пропущенных дат в ежедневной таблице

WITH date_spine AS (
    -- Генерируем все даты от первой записи до вчера
    SELECT explode(sequence(
        (SELECT MIN(event_date) FROM {{ ref('fct_orders') }}),
        CURRENT_DATE() - INTERVAL 1 DAY,
        INTERVAL 1 DAY
    )) AS expected_date
),

actual_dates AS (
    SELECT DISTINCT event_date AS actual_date
    FROM {{ ref('fct_orders') }}
)

SELECT expected_date AS missing_date
FROM date_spine
LEFT JOIN actual_dates ON expected_date = actual_date
WHERE actual_date IS NULL

Тест обнаружит дыры в данных - дни, за которые данные не загружены. Это типичная проблема при инкрементальной загрузке с watermark: если пайплайн не запускался 2 дня, watermark перепрыгнет через эти дни.

Тест для CDC-пайплайнов: дедупликация

-- tests/assert_cdc_no_duplicate_events.sql
-- В CDC-потоке не должно быть дублей одного события

SELECT
    event_id,
    COUNT(*) AS occurrence_count
FROM {{ ref('stg_orders_cdc') }}
GROUP BY event_id
HAVING COUNT(*) > 1

CDC-системы (Debezium, Maxwell, AWS DMS) могут присылать дубли при exactly-once семантике без гарантий. Этот тест ловит их до того, как они попадут в Silver-слой.

Производительность singular тестов

Singular tests запускаются на Spark как обычные SELECT. Для больших таблиц это может быть дорого. Оптимизационные стратегии:

-- Вариант 1: ограничить проверку свежими данными
SELECT *
FROM {{ ref('fct_orders') }}
WHERE event_date >= CURRENT_DATE() - INTERVAL 7 DAYS
  AND discount_amount > total_amount

-- Вариант 2: использовать TABLESAMPLE для приблизительной проверки
-- (не гарантирует 100% покрытие, но быстро)
SELECT *
FROM {{ ref('fct_orders') }} TABLESAMPLE (1 PERCENT)
WHERE total_amount < 0

Для critical финансовых данных используйте полный скан. Для аномалий (отрицательные суммы) часто достаточно 1% выборки: если аномалия системная - она попадёт в выборку.


Data Contracts: строгие гарантии схем

Проблема Schema Drift

Schema Drift - это когда схема таблицы неожиданно меняется без согласования. В СУБД изменение схемы требует явного ALTER TABLE с правами и контролем. В Data Lake инженер бэкенда переименовал колонку в PostgreSQL - CDC обновил файлы в MinIO - Parquet-схема изменилась - Silver-слой сломался, потому что ожидал старое имя.

Без data contracts такие изменения обнаруживаются только при падении downstream-пайплайнов. С data contracts изменение схемы upstream вызывает ошибку при следующем запуске dbt - ещё до того, как некорректные данные попадут в Gold.

Data Contracts в dbt

Начиная с dbt Core 1.5, в dbt появились Data Contracts - возможность формально описать ожидаемую схему модели и включить её принудительную проверку.

Контракт задаётся в schema.yml:

models:
  - name: fct_orders
    description: Gold-витрина заказов с Data Contract

    # Включаем строгую проверку контракта
    config:
      contract:
        enforced: true

    columns:
      - name: order_id
        description: UUID заказа
        data_type: string
        constraints:
          - type: not_null

      - name: user_id
        description: UUID пользователя
        data_type: string
        constraints:
          - type: not_null

      - name: total_amount
        description: Итоговая сумма заказа в рублях
        data_type: decimal(18, 2)
        constraints:
          - type: not_null

      - name: status
        description: Статус заказа
        data_type: string
        constraints:
          - type: not_null

      - name: created_at
        description: Время создания заказа UTC
        data_type: timestamp
        constraints:
          - type: not_null

      - name: event_date
        description: Дата создания (partition key)
        data_type: date
        constraints:
          - type: not_null

Как dbt проверяет контракт

При включённом contract.enforced: true dbt выполняет дополнительный шаг: перед выполнением модели сравнивает объявленные типы колонок с реальными типами в SELECT.

dbt run --select fct_orders

Если SELECT возвращает total_amount как double вместо decimal(18,2):

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

  Contract enforcement requires explicit casting in your SELECT.
  Update the SELECT to: CAST(total_amount AS DECIMAL(18, 2)) AS total_amount

Контракт заставляет писать явные CAST в SELECT-запросах моделей. Это хорошая практика: типы данных становятся документацией прямо в коде, а не скрытым свойством, которое нужно выяснять из метаданных.

Явные типы в модели с контрактом

С включённым контрактом модель должна явно приводить типы:

-- models/marts/fct_orders.sql
{{ config(
    materialized='incremental',
    file_format='delta',
    incremental_strategy='merge',
    unique_key='order_id',
    contract={'enforced': true}
) }}

SELECT
    CAST(o.order_id     AS STRING)          AS order_id,
    CAST(o.user_id      AS STRING)          AS user_id,
    CAST(o.total_amount AS DECIMAL(18, 2))  AS total_amount,
    CAST(o.status       AS STRING)          AS status,
    CAST(o.created_at   AS TIMESTAMP)       AS created_at,
    CAST(date(o.created_at) AS DATE)        AS event_date

FROM {{ ref('int_orders_enriched') }} o
{% if is_incremental() %}
WHERE o.updated_at > (SELECT MAX(updated_at) FROM {{ this }})
{% endif %}

Явные CAST - это не только требование контракта, но и защита от неявного преобразования типов (о котором мы говорили в уроке про ANSI mode).

Ограничения (constraints) в контракте

Контракт поддерживает несколько типов ограничений:

columns:
  - name: order_id
    data_type: string
    constraints:
      - type: not_null        # Обязательное поле

  - name: total_amount
    data_type: decimal(18, 2)
    constraints:
      - type: not_null
      - type: check            # Произвольное SQL-условие
        expression: "total_amount >= 0"
        name: chk_positive_amount

  - name: status
    data_type: string
    constraints:
      - type: not_null

Поддержка check и других ограничений зависит от движка. В Spark 3.x через Delta Lake:

-- Что dbt генерирует для Delta:
ALTER TABLE my_schema.fct_orders
ADD CONSTRAINT chk_positive_amount CHECK (total_amount >= 0)

Для Parquet-таблиц check ограничение логируется, но физически не применяется на уровне движка - оно работает только как документация и валидируется dbt при компиляции.

Versioned Data Contracts

С dbt 1.6 появилась поддержка версионирования моделей - важный инструмент для evolving contracts:

models:
  - name: fct_orders
    latest_version: 2
    versions:
      - v: 1
        defined_in: fct_orders_v1
        # Старый контракт - deprecated, но ещё в production
      - v: 2
        # Текущий контракт
        columns:
          - name: order_id
            data_type: string
          # + новые колонки

Это позволяет постепенно мигрировать потребителей с версии 1 на версию 2 без breaking changes.


Диаграмма: архитектура контроля качества в dbt

Схема показывает полный контур контроля качества данных. Данные проходят несколько уровней защиты: Source Freshness проверяет актуальность ещё на Bronze-уровне, Generic Tests валидируют структурные инварианты на Silver, Singular Tests и Data Contracts защищают Gold-витрины. Все результаты собираются командой dbt build и направляются в CI/CD и систему алертов. Ни один уровень не является лишним: каждый ловит свой класс проблем.


Testing по слоям: разный уровень строгости

Bronze: минимальная валидация

Bronze-слой - это сырые данные «как есть». Здесь жёсткая валидация нежелательна: она будет срабатывать на любые аномалии источника и останавливать ingestion. Цель Bronze - сохранить данные с минимальными изменениями.

# models/staging/sources.yml (Bronze тесты)
sources:
  - name: bronze
    tables:
      - name: raw_orders
        columns:
          - name: _id          # Внутренний ID CDC-события
            tests:
              - not_null        # Единственный hard-fail тест
          - name: order_id
            tests:
              - not_null:
                  severity: warn  # Только предупреждение

Silver: бизнес-правила на очищенных данных

Silver - это очищенные, нормализованные данные. Здесь уже можно применять строгие проверки структурных инвариантов.

# models/staging/schema.yml (Silver тесты)
models:
  - name: stg_orders
    columns:
      - name: order_id
        tests:
          - not_null
          - unique            # Строго: дубли недопустимы после очистки

      - name: user_id
        tests:
          - not_null
          - relationships:
              to: ref('stg_users')
              field: user_id
              severity: warn  # warn: CDC может опережать

      - name: status
        tests:
          - accepted_values:
              values: ['pending', 'confirmed', 'shipped', 'delivered', 'cancelled', 'refunded']

Gold: строгие гарантии для BI и ML

Gold-витрины потребляют BI-инструменты и ML-модели. Здесь максимальная строгость: data contracts, singular tests, severity: error везде.

# models/marts/schema.yml (Gold тесты)
models:
  - name: fct_orders
    config:
      contract:
        enforced: true
    tests:
      - dbt_utils.unique_combination_of_columns:
          combination_of_columns: [order_id]
    columns:
      - name: order_id
        data_type: string
        constraints:
          - type: not_null
        tests:
          - not_null
          - unique

      - name: total_amount
        data_type: decimal(18, 2)
        constraints:
          - type: not_null
        tests:
          - not_null
          - dbt_utils.accepted_range:
              min_value: 0

dbt build: тестирование и трансформации в одном DAG

Разница между dbt run и dbt build

dbt run    # Только выполнение моделей, тесты не запускаются
dbt test   # Только тесты, модели не пересчитываются
dbt build  # Модели + тесты в порядке DAG, с зависимостями

dbt build - правильный способ запускать в production. Он выполняет каждый узел DAG (модель, тест, snapshot, seed) в правильном порядке:

Source Freshness
  → stg_orders (model)
    → not_null_stg_orders_order_id (test)
    → unique_stg_orders_order_id (test)
      → int_orders_enriched (model)           ← ждёт прохождения тестов stg_orders
        → fct_orders (model)
          → assert_discount_not_exceeds_total (test)
          → not_null_fct_orders_order_id (test)

Ключевое свойство dbt build: если тест для stg_orders упал - int_orders_enriched и fct_orders не запустятся. Некачественные данные не попадут вниз по DAG.

CI/CD интеграция

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

on:
  pull_request:
    branches: [main]

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

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

      - name: Check source freshness
        run: dbt source freshness
        env:
          DBT_TARGET: ci
          SPARK_HOST: ${{ secrets.SPARK_HOST }}
          SPARK_TOKEN: ${{ secrets.SPARK_TOKEN }}

      - name: Run dbt build (models + tests)
        run: dbt build --select state:modified+
        # state:modified+ = изменённые модели + все downstream зависимости
        env:
          DBT_TARGET: ci
          SPARK_HOST: ${{ secrets.SPARK_HOST }}
          SPARK_TOKEN: ${{ secrets.SPARK_TOKEN }}

      - name: Upload test results
        if: always()
        uses: actions/upload-artifact@v4
        with:
          name: dbt-test-results
          path: target/run_results.json

--select state:modified+ - это умная фишка dbt: CI запускает только изменённые модели и их downstream зависимости. Это существенно ускоряет CI при большом проекте.


Лабораторная практика: контур защиты данных

Постановка задачи

Мы строим пайплайн для витрины финансовых транзакций. Система должна:

  1. Контролировать свежесть Bronze-данных (SLA: warn через 30 минут, error через 2 часа)
  2. Валидировать структурные инварианты на Silver-слое
  3. Применять бизнес-правила на Gold-уровне
  4. Иметь строгий data contract для Gold-витрины

Структура проекта

models/
├── staging/
│   ├── sources.yml          ← источники + freshness
│   └── stg_transactions.sql
├── intermediate/
│   └── int_transactions_enriched.sql
└── marts/
    ├── schema.yml           ← data contracts + generic tests
    ├── fct_transactions.sql
    └── dim_users.sql

tests/
├── assert_no_negative_amounts.sql
├── assert_discount_valid.sql
└── assert_transaction_balance_reconciliation.sql

Шаг 1: Настройка Source Freshness

# models/staging/sources.yml
version: 2

sources:
  - name: bronze_transactions
    description: CDC-репликация из PostgreSQL, обновляется каждые 15 минут
    database: hive_metastore
    schema: bronze
    loaded_at_field: _ingested_at

    freshness:
      warn_after: {count: 30, period: minutes}
      error_after: {count: 2, period: hours}

    tables:
      - name: raw_transactions
        description: Финансовые транзакции из payments-сервиса
        columns:
          - name: _ingested_at
            description: Timestamp загрузки в Bronze (UTC)
          - name: transaction_id
            description: UUID транзакции
          - name: user_id
            description: UUID пользователя
          - name: amount
            description: Сумма транзакции (может быть отрицательной для возвратов)
          - name: currency
            description: Валюта (ISO 4217)
          - name: status
            description: Статус транзакции

      - name: raw_users
        description: Профили пользователей - обновляются редко
        freshness:
          warn_after: {count: 6, period: hours}
          error_after: {count: 24, period: hours}

Запуск проверки:

dbt source freshness

# Ожидаемый вывод при нормальной работе:
# 14:30:01  Freshness of bronze_transactions.raw_transactions .. [pass in 0.84s]
# 14:30:02  Freshness of bronze_transactions.raw_users ......... [pass in 0.62s]

Для демонстрации WARN - остановим мок-инженер данных на 40 минут и перезапустим:

# Симулируем задержку: обновляем _ingested_at на прошлое
# (в реальности просто ждём или останавливаем пайплайн)
dbt source freshness
# Freshness of bronze_transactions.raw_transactions .. [WARN in 0.91s]
# Elapsed time since last record: 41 minutes (warn threshold: 30 minutes)

Шаг 2: Generic тесты для Silver-слоя

# models/staging/schema.yml
version: 2

models:
  - name: stg_transactions
    description: Очищенные транзакции из Bronze. Одна строка - одна транзакция.

    columns:
      - name: transaction_id
        description: UUID транзакции - первичный ключ
        tests:
          - not_null
          - unique

      - name: user_id
        description: UUID пользователя
        tests:
          - not_null
          - relationships:
              to: ref('stg_users')
              field: user_id
              severity: warn  # CDC может опережать users-таблицу

      - name: amount
        description: Сумма транзакции в копейках (целое число)
        tests:
          - not_null

      - name: currency
        description: Валюта по ISO 4217
        tests:
          - not_null
          - accepted_values:
              values: ['RUB', 'USD', 'EUR', 'CNY']

      - name: status
        description: Статус транзакции
        tests:
          - not_null
          - accepted_values:
              values: ['pending', 'completed', 'failed', 'refunded', 'disputed']

      - name: created_at
        description: Время создания транзакции UTC
        tests:
          - not_null

Шаг 3: Singular тест - аномалии в суммах

-- tests/assert_no_negative_completed_amounts.sql
-- Завершённые транзакции не могут иметь нулевую или отрицательную сумму.
-- Отрицательные суммы допустимы только для статуса 'refunded'.

SELECT
    transaction_id,
    user_id,
    amount,
    status,
    created_at
FROM {{ ref('stg_transactions') }}
WHERE status = 'completed'
  AND amount <= 0
-- tests/assert_refund_not_exceeds_original.sql
-- Сумма возврата не может превышать сумму оригинальной транзакции.
-- Ищем пользователей, чьи возвраты за день больше пополнений.

WITH daily_by_user AS (
    SELECT
        user_id,
        event_date,
        SUM(CASE WHEN status = 'completed' THEN amount ELSE 0 END) AS total_credits,
        SUM(CASE WHEN status = 'refunded'  THEN ABS(amount) ELSE 0 END) AS total_refunds
    FROM {{ ref('fct_transactions') }}
    GROUP BY user_id, event_date
)

SELECT
    user_id,
    event_date,
    total_credits,
    total_refunds,
    (total_refunds - total_credits) AS overrefund
FROM daily_by_user
WHERE total_refunds > total_credits

Шаг 4: Data Contract для Gold-витрины

# models/marts/schema.yml
version: 2

models:
  - name: fct_transactions
    description: >
      Gold-витрина финансовых транзакций.
      Строгий Data Contract - изменения схемы требуют согласования.

    config:
      contract:
        enforced: true

    columns:
      - name: transaction_id
        description: UUID транзакции
        data_type: string
        constraints:
          - type: not_null
        tests:
          - not_null
          - unique

      - name: user_id
        description: UUID пользователя
        data_type: string
        constraints:
          - type: not_null
        tests:
          - not_null

      - name: amount_rub
        description: Сумма транзакции в рублях (с копейками)
        data_type: decimal(18, 2)
        constraints:
          - type: not_null
        tests:
          - not_null
          - dbt_utils.accepted_range:
              min_value: -1000000   # Максимальный возврат 1 млн руб
              max_value: 10000000   # Максимальная транзакция 10 млн руб

      - name: status
        description: Финальный статус транзакции
        data_type: string
        constraints:
          - type: not_null
        tests:
          - not_null
          - accepted_values:
              values: ['completed', 'refunded', 'failed', 'disputed']

      - name: created_at
        description: Время создания транзакции UTC
        data_type: timestamp
        constraints:
          - type: not_null
        tests:
          - not_null

      - name: event_date
        description: Дата создания (partition key)
        data_type: date
        constraints:
          - type: not_null
        tests:
          - not_null

Модель с явными типами:

-- models/marts/fct_transactions.sql
{{ config(
    materialized='incremental',
    file_format='delta',
    incremental_strategy='merge',
    unique_key='transaction_id',
    partition_by={'field': 'event_date', 'data_type': 'date'},
    contract={'enforced': true}
) }}

SELECT
    CAST(t.transaction_id           AS STRING)          AS transaction_id,
    CAST(t.user_id                  AS STRING)          AS user_id,
    CAST(t.amount / 100.0           AS DECIMAL(18, 2))  AS amount_rub,
    CAST(t.status                   AS STRING)          AS status,
    CAST(t.created_at               AS TIMESTAMP)       AS created_at,
    CAST(date(t.created_at)         AS DATE)            AS event_date,
    CAST(current_timestamp()        AS TIMESTAMP)       AS _loaded_at

FROM {{ ref('int_transactions_enriched') }} t

{% if is_incremental() %}
WHERE t.updated_at > (SELECT MAX(_loaded_at) FROM {{ this }})
{% endif %}

Шаг 5: Запуск полного контура качества

# 1. Проверить свежесть источников
dbt source freshness

# 2. Запустить все модели с тестами (в порядке DAG)
dbt build

# Ожидаемый вывод:
# 14:30:01  Running 8 nodes in thread pool with 4 threads
# 14:30:02  1 of 8 START source freshness check ... [RUN]
# 14:30:02  1 of 8 PASS source freshness check ... [pass in 0.84s]
# 14:30:03  2 of 8 START model stg_transactions ... [RUN]
# 14:30:15  2 of 8 OK created model stg_transactions ... [OK in 12.4s]
# 14:30:15  3 of 8 START test not_null_stg_transactions_transaction_id ... [RUN]
# 14:30:16  3 of 8 PASS test not_null_stg_transactions_transaction_id ... [PASS in 0.91s]
# 14:30:16  4 of 8 START test unique_stg_transactions_transaction_id ... [RUN]
# 14:30:17  4 of 8 PASS test unique_stg_transactions_transaction_id ... [PASS in 1.12s]
# ...
# 14:31:45  8 of 8 START test assert_refund_not_exceeds_original ... [RUN]
# 14:31:50  8 of 8 PASS test assert_refund_not_exceeds_original ... [PASS in 4.91s]
#
# Finished running 3 models, 5 tests in 1 minutes 50.34 seconds.
# PASS=8 WARN=0 ERROR=0 SKIP=0 TOTAL=8

# 3. Посмотреть скомпилированный SQL тестов
ls target/compiled/my_project/tests/
cat target/compiled/my_project/tests/assert_refund_not_exceeds_original.sql

Шаг 6: Анализ упавшего теста в Spark UI

При падении теста dbt показывает:

14:31:50  3 of 8 FAIL 2 unique_fct_transactions_transaction_id ... [FAIL 2 in 1.12s]

Число 2 - количество нарушителей. Для детального анализа:

# Включить store_failures и перезапустить
dbt test --select fct_transactions --store-failures

# Посмотреть нарушителей
dbt run-operation run_query --args '{query: "SELECT * FROM dbt_test__audit.unique_fct_transactions_transaction_id LIMIT 10"}'

В Spark UI (вкладка SQL) видим план выполнения теста: это обычный SELECT с GROUP BY и HAVING COUNT(*) > 1. При большой таблице Spark выполняет его через HashAggregate + ShuffleExchange - по сути, распределённый подсчёт дублей по всему кластеру.


Best Practices: тестирование на больших данных

Управление стоимостью тестов

Тесты на Spark - это полноценные Spark-джобы. Каждый dbt test запускает SELECT по всей (или части) таблицы. При неосторожном подходе стоимость тестирования может превысить стоимость самих трансформаций.

Стратегии оптимизации:

# 1. Используйте where для ограничения объёма
- name: order_id
  tests:
    - unique:
        where: "event_date >= current_date() - interval 7 days"

# 2. Настройте severity правильно:
# error - только для критичных бизнес-инвариантов (финансы, ключи)
# warn  - для мягких проверок (relationships при асинхронной загрузке)

# 3. Разделяйте тесты на быстрые и медленные:
# dbt test --select "tag:fast"
# dbt test --select "tag:slow"  (запускать реже)
# Тегирование тестов
models:
  - name: fct_orders
    columns:
      - name: order_id
        tests:
          - not_null:
              meta:
                tags: ['fast', 'critical']
          - unique:
              meta:
                tags: ['fast', 'critical']

tests:
  - name: assert_transaction_balance_reconciliation
    meta:
      tags: ['slow', 'financial']

Anti-patterns в тестировании

Тестировать только not_null - это минималистичный подход. Пропускаются: дубли (unique), выбросы (accepted_range), нарушения бизнес-логики (singular tests).

Игнорировать freshness - самый частый пропуск. Свежесть данных - это первое, что нужно проверять в production-пайплайне.

Тестировать только Gold-слой - слишком поздно. Проблема, обнаруженная на Gold, уже «заразила» несколько слоёв. Тестируйте на каждом слое с разной строгостью.

Все тесты через dbt run + dbt test вместо dbt build. Это позволяет некачественным данным попасть в downstream-модели. Используйте dbt build.

Отсутствие store_failures на критичных тестах. При падении теста вы не знаете, какие конкретно строки проблемны. store_failures: true решает это.

Дорогие relationships тесты на больших таблицах - NOT IN (subquery) на 100 GB + 50 GB = shuffle двух таблиц. Добавляйте where или снижайте до severity: warn.


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

Условие задачи

Дан dbt-проект с тремя таблицами:

  • stg_products - справочник товаров
  • fct_sales - продажи (100 GB, партиционированы по sale_date)
  • dim_customers - профили клиентов

В данных обнаружены проблемы: периодически появляются дубликаты по sale_id, иногда проскакивают продажи с отрицательной маржой, а источники обновляются с задержкой.

Задание 1: Полное покрытие generic тестами

Напишите schema.yml для всех трёх таблиц:

  1. Покройте все бизнес-ключи тестами not_null + unique.
  2. Добавьте relationships там, где есть внешние ссылки (customer_id, product_id).
  3. Добавьте accepted_values для всех статусных полей.
  4. Для полей amount и margin добавьте dbt_utils.accepted_range с адекватными границами.

Задание 2: Singular тест на сверку баланса

Напишите тест assert_margin_consistency.sql:

  • Маржинальность по каждому товару должна совпадать в fct_sales и в агрегированной таблице agg_product_margins
  • Допустимое расхождение - не более 0.01 рубля (ошибки округления)

Задание 3: Source Freshness

Настройте freshness для источников:

  • raw_sales: warn через 15 минут, error через 1 час (OLTP-система обновляется часто)
  • raw_products: warn через 6 часов, error через 24 часа (справочник меняется редко)
  • raw_customers: warn через 1 час, error через 6 часов

Задание 4: Data Contract для Gold-витрины

Включите contract: enforced: true для fct_sales. Укажите явные data_type для каждой колонки. Обновите SQL-модель с явными CAST.

Что сдавать

  1. Файлы schema.yml, sources.yml, tests/assert_margin_consistency.sql
  2. Лог успешного выполнения dbt build (вывод из терминала)
  3. Содержимое target/compiled/my_project/tests/assert_margin_consistency.sql - скомпилированный SQL без Jinja

Это задание проверяет, что вы умеете не только настроить тесты, но и понимаете, в какой SQL они компилируются и как Spark их выполняет.