dbt Tests и Source Freshness: generic tests, singular tests, data contracts
dbt Tests и Source Freshness: generic tests, singular tests, data contracts
Концепция 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- колонка не содержит NULLunique- все значения уникальны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 при большом проекте.
Лабораторная практика: контур защиты данных¶
Постановка задачи¶
Мы строим пайплайн для витрины финансовых транзакций. Система должна:
- Контролировать свежесть Bronze-данных (SLA: warn через 30 минут, error через 2 часа)
- Валидировать структурные инварианты на Silver-слое
- Применять бизнес-правила на Gold-уровне
- Иметь строгий 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 для всех трёх таблиц:
- Покройте все бизнес-ключи тестами
not_null+unique. - Добавьте
relationshipsтам, где есть внешние ссылки (customer_id, product_id). - Добавьте
accepted_valuesдля всех статусных полей. - Для полей
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.
Что сдавать¶
- Файлы
schema.yml,sources.yml,tests/assert_margin_consistency.sql - Лог успешного выполнения
dbt build(вывод из терминала) - Содержимое
target/compiled/my_project/tests/assert_margin_consistency.sql- скомпилированный SQL без Jinja
Это задание проверяет, что вы умеете не только настроить тесты, но и понимаете, в какой SQL они компилируются и как Spark их выполняет.