dbt macros для Spark: partition_by, clustered_by, location_root и submit_timeout
dbt macros для Spark: partition_by, clustered_by, location_root и submit_timeout
Почему Spark требует физических конфигураций¶
Когда дата-инженер переходит от привычного SQL к dbt + Spark, он обнаруживает неприятный факт: написать правильный SELECT недостаточно. Запрос может возвращать корректные данные, но при этом создавать миллионы маленьких файлов, делать каждый join в десятки раз медленнее обычного, или хранить всё в одной огромной директории, которую Metastore сканирует минутами.
Проблема в том, что Spark - это не традиционная СУБД со встроенным оптимизатором хранения. Spark - это вычислительный движок, который честно делает то, что ему сказали. Если вы не задали партиционирование, Spark запишет все данные в одну директорию. Если не задали bucketing, каждый join будет сопровождаться полным shuffle. Если не ограничили timeout, долгая инициализация кластера тихо провалится без понятной ошибки.
Традиционные СУБД (PostgreSQL, ClickHouse) имеют встроенную статистику и автоматически выбирают физическое хранение. Spark переносит ответственность за физическую организацию данных на инженера. Это мощь - но она требует явных решений.
Проблема "хорошего SQL, плохих данных"¶
Рассмотрим конкретный пример. Витрина fct_transactions обрабатывает 100 GB данных в сутки за 4 года - итого 146 TB. Простой dbt-model без физических настроек создаст:
-- dbt создаст без физических настроек:
-- 1. Одну плоскую директорию с ~50 000 файлов по 3 MB
-- 2. Каждый запрос по дате = full scan 146 TB
-- 3. JOIN с user_dim = shuffle всех 146 TB
-- 4. Нет возможности параллельного чтения по диапазону дат
С правильными физическими настройками:
-- С partition_by + clustered_by:
-- 1. 1 460 партиций по дате, каждая ~100 GB
-- 2. Запрос по дате = scan только нужной партиции (100 GB вместо 146 TB)
-- 3. JOIN с user_dim = Shuffle-Free, файлы уже отсортированы по user_id
-- 4. Metastore знает структуру, pruning работает автоматически
Разница в производительности: от часов к минутам. Именно для управления этой физической организацией существуют Spark-специфичные конфигурации в dbt.
Место физических конфигураций в dbt¶
dbt предоставляет несколько уровней конфигурации, которые применяются в порядке приоритета (от низшего к высшему):
dbt_project.yml (глобальные defaults)
└── schema.yml (конфигурация директории)
└── {{ config() }} в модели (конфигурация файла)
Spark-специфичные параметры (partition_by, clustered_by, location_root, file_format, submit_timeout) передаются через этот же механизм и транслируются адаптером в DDL и настройки сессии. Это ключевое отличие от обычного SQL: вы не пишете PARTITIONED BY явно - вы объявляете желаемое физическое устройство, а dbt генерирует нужный DDL.
Jinja как язык dbt: основы для практики¶
Прежде чем разбирать конкретные макросы, важно понять, что dbt - это не просто SQL-препроцессор. dbt использует Jinja2 как полноценный язык шаблонов, который выполняется на стороне клиента (на машине, где запускается dbt) до отправки SQL на Spark.
Три типа Jinja-блоков¶
{# Это комментарий - не попадает в результирующий SQL #}
{{ expression }} {# Вычислить и вставить значение в текст #}
{% statement %} {# Управляющая конструкция: if, for, set, call #}
{% endstatement %}
Переменные и контекст¶
В dbt-шаблонах доступны специальные переменные:
{{ this }} {# Ссылка на текущую модель: database.schema.model_name #}
{{ model }} {# Объект текущей модели: имя, путь, конфиг #}
{{ target }} {# Текущий target: name, schema, database, type #}
{{ env_var('KEY') }} {# Переменная окружения, обязательная #}
{{ env_var('KEY', 'default') }} {# С default-значением #}
{{ var('my_var') }} {# Переменная из dbt_project.yml или --vars CLI #}
Макросы: переиспользуемые блоки логики¶
Макрос в dbt - это именованный Jinja-шаблон, который принимает аргументы и возвращает строку:
{# macros/utils.sql #}
{% macro my_macro(arg1, arg2='default') %}
{# тело макроса - любая Jinja-логика #}
SELECT {{ arg1 }}, '{{ arg2 }}'
{% endmacro %}
Макросы вызываются из моделей или других макросов:
-- models/my_model.sql
{{ my_macro('column_name') }}
-- или через call, если макрос не возвращает SQL:
{% call(result) my_macro('arg') %}{% endcall %}
Жизненный цикл макроса¶
Понимание того, когда выполняется Jinja - критично для отладки:
dbt run
│
├─ 1. Jinja compile phase (локально, до Spark)
│ Все {{ }}, {% %} раскрываются в plain SQL
│ Макросы выполняются здесь
│ env_var(), var() читаются здесь
│
├─ 2. SQL отправляется на кластер через адаптер
│
├─ 3. Spark выполняет SQL
│ CREATE TABLE, INSERT, MERGE - здесь
│
└─ 4. dbt читает метаданные результата
Вывод: макросы не выполняются на Spark. Они генерируют SQL, который потом отправляется на Spark. Это объясняет, почему нельзя использовать Spark UDF внутри макроса - к моменту выполнения UDF Jinja уже завершила работу.
Отладка макросов: dbt compile¶
Самый мощный инструмент отладки - dbt compile. Он раскрывает все Jinja-шаблоны в готовый SQL без отправки на Spark:
# Скомпилировать конкретную модель
dbt compile --select fct_transactions
# Результат появится в target/compiled/my_project/models/fct_transactions.sql
cat target/compiled/my_project/models/fct_transactions.sql
Всегда используйте dbt compile при разработке макросов - это позволяет видеть генерируемый DDL до фактического выполнения.
Макрос partition_by: физическое партиционирование таблиц¶
Партиционирование - самый важный механизм оптимизации хранения в Spark. Правильно выбранный ключ партиции превращает полный скан 146 TB в точечный запрос 100 GB.
Что такое партиционирование в Hive-стиле¶
Spark (через HMS - Hive Metastore) поддерживает партиционирование таблиц по значениям столбцов. Физически это означает:
/warehouse/fct_transactions/
├── event_date=2024-01-01/ ← partition directory
│ ├── part-00000.parquet
│ └── part-00001.parquet
├── event_date=2024-01-02/
│ ├── part-00000.parquet
│ └── part-00001.parquet
└── ...
Когда запрос содержит WHERE event_date = '2024-01-01', Spark читает только соответствующую поддиректорию. Это называется partition pruning - механизм, при котором Metastore исключает нерелевантные партиции ещё до чтения данных.
Синтаксис partition_by в dbt-spark¶
В dbt-spark параметр partition_by передаётся через config():
-- models/fct_transactions.sql
{{ config(
materialized='incremental',
file_format='parquet',
partition_by={
"field": "event_date",
"data_type": "date"
}
) }}
SELECT
transaction_id,
user_id,
amount,
event_date
FROM {{ ref('int_transactions_enriched') }}
Адаптер транслирует это в DDL:
-- Что dbt генерирует для Parquet/Hive-формата:
CREATE TABLE my_schema.fct_transactions
USING PARQUET
PARTITIONED BY (event_date)
AS SELECT ...
-- При инкрементальном добавлении:
INSERT INTO my_schema.fct_transactions
PARTITION (event_date)
SELECT ...
Типы данных partition_by¶
Выбор типа данных для партиции влияет на поведение:
| data_type | Пример значения | Количество партиций | Рекомендация |
|---|---|---|---|
date |
2024-01-15 |
365/год | Хорошо для суточных данных |
timestamp |
2024-01-15 12:00:00 |
Огромное | Плохо, слишком мелко |
int |
20240115 (yyyyMMdd) |
365/год | Допустимо, но хуже date |
string |
"2024-01-15" |
365/год | Только если нет другого варианта |
month (Iceberg) |
2024-01 |
12/год | Отлично для крупных витрин |
Никогда не партиционируйте по timestamp без трансформации - это создаёт уникальную партицию для каждой записи и убивает Metastore.
Многоуровневое партиционирование¶
dbt-spark поддерживает составные ключи партиционирования:
{{ config(
materialized='table',
file_format='parquet',
partition_by=[
{"field": "country_code", "data_type": "string"},
{"field": "event_date", "data_type": "date"}
]
) }}
Физическая структура директорий:
/warehouse/fct_transactions/
├── country_code=RU/
│ ├── event_date=2024-01-01/
│ │ └── part-00000.parquet
│ └── event_date=2024-01-02/
│ └── part-00001.parquet
├── country_code=DE/
│ └── event_date=2024-01-01/
│ └── part-00002.parquet
└── ...
Запрос WHERE country_code = 'RU' AND event_date = '2024-01-01' читает только одну маленькую поддиректорию.
Важное правило: порядок полей в partition_by важен. Первое поле - верхний уровень иерархии. Запросы без фильтра по первому полю не получают pruning по второму. Ставьте поле с наименьшей кардинальностью первым (страна, год), а с наибольшей - последним (дата).
partition_by для Apache Iceberg: скрытое партиционирование¶
Когда file_format='iceberg', поведение partition_by меняется принципиально. Iceberg поддерживает скрытое партиционирование (hidden partitioning) - трансформации прямо в конфигурации:
{{ config(
materialized='incremental',
file_format='iceberg',
partition_by={
"field": "created_at",
"data_type": "timestamp",
"granularity": "day" -- скрытая трансформация!
}
) }}
Адаптер генерирует Iceberg DDL:
CREATE TABLE my_catalog.my_schema.fct_transactions
USING iceberg
PARTITIONED BY (days(created_at))
Пользователь запрашивает WHERE created_at > '2024-01-01' - Iceberg автоматически находит нужные партиции, хотя в данных нет отдельного столбца event_date. Никаких date_trunc() в запросах.
Поддерживаемые granularity для Iceberg:
-- В config():
"granularity": "year" -- PARTITIONED BY (years(ts))
"granularity": "month" -- PARTITIONED BY (months(ts))
"granularity": "day" -- PARTITIONED BY (days(ts))
"granularity": "hour" -- PARTITIONED BY (hours(ts))
Для нечастотных трансформаций (bucket, truncate) используют явный tblproperties:
{{ config(
materialized='table',
file_format='iceberg',
tblproperties={
"write.distribution-mode": "hash"
},
partition_by={
"field": "user_id",
"data_type": "int",
"transform": "bucket[32]" -- bucket(32, user_id)
}
) }}
Как проверить партиционирование¶
После dbt run - проверить через SQL:
-- Для Hive/Parquet - показать партиции
SHOW PARTITIONS my_schema.fct_transactions;
-- Для Iceberg - через metadata tables
SELECT partition, record_count, file_count
FROM my_catalog.my_schema.fct_transactions.partitions;
-- Проверить partition pruning в EXPLAIN
EXPLAIN SELECT * FROM my_schema.fct_transactions
WHERE event_date = '2024-01-01';
-- В плане должно быть "PartitionFilters: [isnotnull(event_date#X), (event_date#X = 18276)]"
Макрос clustered_by и buckets: bucketing и Shuffle-Free Join¶
Если partition_by решает проблему partition pruning при фильтрации, то clustered_by решает проблему shuffle при join. Это два ортогональных механизма оптимизации, и для больших витрин нужны оба.
Проблема shuffle при join¶
Когда Spark выполняет join двух таблиц, он должен убедиться, что строки с одинаковым ключом join находятся на одном executore. Стандартный способ - shuffle: перераспределить данные обеих таблиц по ключу join через сеть.
Стоимость shuffle огромна:
Без bucketing (стандартный join):
1. Serialize 100 GB fct_transactions → память + сеть
2. Shuffle по user_id → 100 GB через network
3. Deserialize на executors → IO
4. Serialize 10 GB user_dim → память + сеть
5. Shuffle по user_id → 10 GB через network
6. SortMergeJoin → CPU
Итого: 2× shuffle overhead
С bucketing данные уже физически распределены по ключу join. Spark может выполнить join без shuffle:
С bucketing (Shuffle-Free Join):
1. Читаем fct_transactions bucket 0 → только нужные данные
2. Читаем user_dim bucket 0 → только нужные данные
3. BucketedMergeJoin bucket by bucket → CPU без сети
Итого: 0 shuffle
Синтаксис clustered_by в dbt-spark¶
{{ config(
materialized='table',
file_format='parquet',
clustered_by=['user_id'],
buckets=256
) }}
SELECT
transaction_id,
user_id,
amount
FROM {{ ref('int_transactions_enriched') }}
Адаптер генерирует DDL с CLUSTERED BY:
CREATE TABLE my_schema.fct_transactions
USING PARQUET
CLUSTERED BY (user_id) INTO 256 BUCKETS
AS SELECT ...
Как работает bucketing¶
Spark вычисляет номер bucket для каждой строки:
bucket_id = hash(user_id) % num_buckets
Файлы записываются с суффиксом, указывающим номер bucket:
/warehouse/fct_transactions/
├── part-00000-a1b2c3.parquet ← bucket 0
├── part-00001-d4e5f6.parquet ← bucket 1
├── ...
└── part-00255-x7y8z9.parquet ← bucket 255
Строки с одинаковым user_id всегда попадают в один и тот же bucket - это гарантия хеш-функции. Поэтому Spark может joинить bucket 0 с bucket 0, bucket 1 с bucket 1 - без перемешивания данных.
Выбор количества buckets¶
Количество buckets - критичный параметр. Нет универсальной формулы, но есть эмпирические правила:
target_bucket_size = 128-256 MB
num_buckets = total_table_size_MB / target_bucket_size
# Примеры:
# 100 GB таблица: 100_000 / 128 = 781 → округлить до 512 (степень двойки)
# 10 GB таблица: 10_000 / 128 = 78 → округлить до 64 или 128
# 1 TB таблица: 1_000_000 / 256 = 3906 → округлить до 4096
Правило степени двойки: количество buckets должно быть степенью двойки (64, 128, 256, 512...). Это позволяет Spark автоматически выполнять bucket coalescing - объединение совместимых таблиц с разным числом buckets (если одна таблица имеет 256 buckets, а другая 512 - Spark сможет выровнять их без shuffle).
Условия для Shuffle-Free Join¶
Shuffle-Free Join работает только при выполнении всех условий:
# Оба условия должны быть истинны:
# 1. Обе таблицы забакетированы по одному ключу
# 2. Число buckets совместимо (равно или кратно)
# Проверить в Spark UI → SQL → план → должно быть:
# "BucketedMergeJoin" вместо "SortMergeJoin + Exchange"
Если одна из таблиц не забакетирована - обычный shuffle. Поэтому имеет смысл бакетировать как основную витрину, так и справочники, если они активно используются в join.
Совмещение partition_by и clustered_by¶
Оба параметра можно комбинировать - это стандартная практика для больших витрин:
{{ config(
materialized='incremental',
file_format='parquet',
incremental_strategy='insert_overwrite',
partition_by={
"field": "event_date",
"data_type": "date"
},
clustered_by=['user_id'],
buckets=256
) }}
Физическая структура:
/warehouse/fct_transactions/
├── event_date=2024-01-01/ ← partition pruning по дате
│ ├── part-00000.parquet ← bucket 0: user_id с hash%256=0
│ ├── part-00001.parquet ← bucket 1
│ └── ...
└── event_date=2024-01-02/
├── part-00000.parquet ← те же самые buckets внутри партиции
└── ...
Результат: запрос WHERE event_date = '2024-01-01' JOIN user_dim USING (user_id) - pruning по дате + Shuffle-Free Join по user_id.
Ограничения bucketing в Iceberg¶
Важное ограничение: классический Hive-bucketing (CLUSTERED BY ... INTO N BUCKETS) - это концепция Parquet/ORC-таблиц в Hive-стиле. Apache Iceberg имеет собственную реализацию bucket-распределения через скрытое партиционирование:
-- Для Iceberg используйте partition_by с transform bucket:
{{ config(
materialized='table',
file_format='iceberg',
partition_by={
"field": "user_id",
"data_type": "bigint",
"transform": "bucket[256]"
}
) }}
Iceberg bucket-transform - более мощный механизм: он совместим между Spark, Flink, Trino и поддерживает Partition Evolution.
Параметр location_root: внешние таблицы и изоляция окружений¶
В production-системах данные часто хранятся не в стандартной директории Metastore, а во внешнем хранилище: S3, HDFS, MinIO. Параметр location_root позволяет dbt создавать external tables с явным путём хранения.
Managed vs External Tables¶
Разница между двумя типами таблиц:
-- Managed table (без location):
-- Данные хранятся в /user/hive/warehouse/
-- DROP TABLE удаляет данные
CREATE TABLE my_schema.fct_transactions
USING PARQUET
AS SELECT ...
-- External table (с location):
-- Данные хранятся по указанному пути
-- DROP TABLE удаляет только метаданные, данные остаются
CREATE TABLE my_schema.fct_transactions
USING PARQUET
LOCATION 's3a://data-lake/gold/fct_transactions'
AS SELECT ...
В production почти всегда используют external tables - это защищает данные от случайного DROP TABLE и позволяет работать с ними из разных инструментов (Spark, Trino, Flink) по одному физическому пути.
Синтаксис location_root в dbt-spark¶
{{ config(
materialized='table',
file_format='parquet',
location_root='s3a://data-lake/gold'
) }}
SELECT * FROM {{ ref('int_transactions_enriched') }}
dbt автоматически дополняет путь именем модели:
-- Итоговый LOCATION:
-- s3a://data-lake/gold/fct_transactions
CREATE TABLE my_schema.fct_transactions
USING PARQUET
LOCATION 's3a://data-lake/gold/fct_transactions'
AS SELECT ...
Динамические пути через env_var¶
Жёсткий путь в config() - антипаттерн. Нельзя хранить разные пути для dev/prod в одном конфиг-блоке. Решение - env_var():
{{ config(
materialized='incremental',
file_format='parquet',
location_root=env_var('DBT_LOCATION_ROOT', 's3a://dev-lake/gold')
) }}
Переменные окружения по среде:
# .env для локальной разработки
DBT_LOCATION_ROOT=s3a://dev-lake/gold
# GitHub Actions secrets для CI
DBT_LOCATION_ROOT=s3a://staging-lake/gold
# Kubernetes Secret для production
DBT_LOCATION_ROOT=s3a://prod-lake/gold
Это гарантирует, что dev-запуски не перезаписывают production-данные - физически разные пути в S3.
Изоляция dev/prod через target¶
Альтернативный подход - использовать target.name (имя активного профиля):
{# macros/get_location_root.sql #}
{% macro get_location_root(layer='gold') %}
{%- if target.name == 'prod' -%}
s3a://prod-lake/{{ layer }}
{%- elif target.name == 'staging' -%}
s3a://staging-lake/{{ layer }}
{%- else -%}
s3a://dev-lake/{{ target.schema }}/{{ layer }}
{%- endif -%}
{% endmacro %}
Использование в модели:
{{ config(
materialized='table',
file_format='parquet',
location_root=get_location_root('gold')
) }}
При dbt run --target prod → s3a://prod-lake/gold/fct_transactions.
При dbt run --target dev → s3a://dev-lake/dev_ivan/gold/fct_transactions.
Включение target.schema в dev-путь позволяет разным разработчикам работать в изолированных пространствах без конфликтов.
Макрос generate_schema_name для изоляции схем¶
В дополнение к физическому пути, dbt позволяет переопределить логику генерации имён схем. Это стандартный макрос-хук:
{# 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 -%}
{%- if target.name == 'prod' -%}
{{ custom_schema_name | trim }}
{%- else -%}
{{ default_schema }}_{{ custom_schema_name | trim }}
{%- endif -%}
{%- endif -%}
{%- endmacro %}
В production: схема gold. В dev: схема dev_ivan_gold. Данные изолированы.
location_root для Iceberg и REST-каталогов¶
При использовании Iceberg с REST-каталогом (Nessie, Polaris, Unity Catalog) location_root работает по-другому. REST-каталог сам управляет расположением данных через warehouse настройку:
# profiles.yml для Iceberg + REST Catalog
my_iceberg_project:
target: dev
outputs:
dev:
type: spark
method: session
schema: dev_ivan
config:
spark.sql.catalog.my_catalog: org.apache.iceberg.spark.SparkCatalog
spark.sql.catalog.my_catalog.type: rest
spark.sql.catalog.my_catalog.uri: http://nessie:19120/api/v1
spark.sql.catalog.my_catalog.warehouse: s3a://data-lake/warehouse
В этом случае location_root в config() игнорируется - расположение данных определяется каталогом. location_root актуален для Hadoop Catalog и для Parquet/ORC-таблиц в Hive-стиле.
submit_timeout и конфигурация Spark-сессии¶
Одна из типичных проблем при работе с dbt + Spark - таймауты при инициализации кластера. Когда кластер масштабируется с нуля (cold start), первый dbt-запрос может ждать несколько минут. По умолчанию dbt обрывает соединение через 30 секунд.
Проблема cold start¶
[DEV] dbt run --select fct_transactions
│
├─ Connecting to Spark Thrift Server...
├─ Submitting query...
│ (Spark кластер начинает запуск executors)
├─ Waiting for response...
│ (30 секунд проходит)
│
└─ ERROR: Connection timeout after 30s
Runtime Error in model fct_transactions
submit_timeout в profiles.yml¶
Параметр submit_timeout (для http-метода) или connect_timeout (для thrift) задаёт максимальное время ожидания ответа:
# profiles.yml
my_project:
target: prod
outputs:
prod:
type: spark
method: http
host: databricks-workspace.azuredatabricks.net
token: "{{ env_var('DATABRICKS_TOKEN') }}"
endpoint: /sql/1.0/warehouses/abc123
schema: gold
threads: 8
submit_timeout: 600 # 10 минут для cold start
connect_timeout: 60 # 1 минута для TCP handshake
connect_retries: 3 # 3 попытки подключения
Для thrift-метода:
prod:
type: spark
method: thrift
host: spark-master.internal
port: 10000
schema: gold
threads: 4
connect_timeout: 120 # ждём 2 минуты TCP-соединения
connect_retries: 5 # 5 попыток перед ошибкой
server_side_parameters:
"spark.sql.shuffle.partitions": "400"
server_side_parameters: конфигурация Spark через dbt¶
server_side_parameters - это словарь Spark-конфигураций, которые dbt автоматически выполняет через SET команды при установке соединения:
# profiles.yml
my_project:
outputs:
prod:
type: spark
method: thrift
host: spark-thrift.internal
schema: gold
server_side_parameters:
# Производительность
"spark.sql.shuffle.partitions": "400"
"spark.sql.adaptive.enabled": "true"
"spark.sql.adaptive.coalescePartitions.enabled": "true"
"spark.sql.adaptive.advisoryPartitionSizeInBytes": "134217728" # 128 MB
# Запись данных
"spark.sql.parquet.compression.codec": "snappy"
"spark.sql.parquet.mergeSchema": "false"
# Delta Lake
"spark.databricks.delta.optimizeWrite.enabled": "true"
"spark.databricks.delta.autoCompact.enabled": "true"
# Iceberg
"spark.sql.extensions": "org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions"
"spark.sql.defaultCatalog": "my_catalog"
При запуске dbt эквивалентно выполнению:
SET spark.sql.shuffle.partitions = 400;
SET spark.sql.adaptive.enabled = true;
-- ... и т.д.
post-hook для Spark-конфигураций на уровне модели¶
Иногда нужно применить конфигурации только для конкретной модели, не для всей сессии. Для этого используются pre_hook и post_hook:
{{ config(
materialized='incremental',
file_format='parquet',
pre_hook=[
"SET spark.sql.shuffle.partitions = 800",
"SET spark.sql.adaptive.skewJoin.enabled = true"
],
post_hook=[
"SET spark.sql.shuffle.partitions = 200",
"ANALYZE TABLE {{ this }} COMPUTE STATISTICS FOR ALL COLUMNS"
]
) }}
SELECT * FROM {{ ref('int_transactions_enriched') }}
pre_hook выполняется до основного SQL модели. post_hook - после. Это позволяет временно изменить конфигурацию Spark под конкретную тяжёлую операцию, а потом вернуть к дефолтам.
Типичный use case pre_hook - настройка shuffle partitions под размер конкретной таблицы:
{{ config(
pre_hook="SET spark.sql.shuffle.partitions = {{ (model_size_gb * 8) | int }}"
) }}
Где model_size_gb - переменная из dbt_project.yml или --vars.
Макрос для динамического расчёта shuffle partitions¶
Создадим макрос, который рассчитывает оптимальное число shuffle partitions на основе размера таблицы:
{# macros/spark_config.sql #}
{% macro optimal_shuffle_partitions(
estimated_size_gb,
target_partition_mb=128,
min_partitions=50,
max_partitions=2000
) %}
{%- set raw = (estimated_size_gb * 1024 / target_partition_mb) | int -%}
{%- set result = [min_partitions, [raw, max_partitions] | min] | max -%}
{{ result }}
{% endmacro %}
Использование в pre_hook:
{{ config(
pre_hook="SET spark.sql.shuffle.partitions = {{ optimal_shuffle_partitions(100) }}"
) }}
-- SET spark.sql.shuffle.partitions = 800
Глобальные конфигурации в dbt_project.yml¶
Прописывать config() в каждом файле модели - неудобно. dbt позволяет задать defaults глобально или на уровне директории в dbt_project.yml:
Структура dbt_project.yml с Spark-конфигами¶
# dbt_project.yml
name: my_data_platform
version: '1.0.0'
config-version: 2
profile: my_project
# Пути к компонентам
model-paths: ["models"]
macro-paths: ["macros"]
test-paths: ["tests"]
# Глобальные defaults для моделей
models:
my_data_platform:
# Всё в проекте - incremental Parquet по умолчанию
+materialized: incremental
+file_format: parquet
+incremental_strategy: insert_overwrite
# Слой staging: view для экономии места
staging:
+materialized: view
# Слой intermediate: ephemeral (без физической таблицы)
intermediate:
+materialized: ephemeral
# Слой marts: полная оптимизация
marts:
+materialized: incremental
+file_format: parquet
+incremental_strategy: insert_overwrite
+location_root: "{{ env_var('DBT_LOCATION_ROOT', 's3a://dev-lake/gold') }}"
# Финансовые витрины - отдельные настройки
finance:
+partition_by:
field: report_date
data_type: date
+clustered_by:
- company_id
+buckets: 128
Теперь модель в директории models/marts/finance/fct_revenue.sql автоматически получает:
materialized='incremental'file_format='parquet'partition_by={'field': 'report_date', 'data_type': 'date'}clustered_by=['company_id']buckets=128location_rootиз переменной окружения
Модель может переопределить любой параметр через свой config() - это всегда приоритетнее глобального.
Переменные через --vars и var()¶
Иногда нужно передавать параметры не через файлы конфигурации, а через CLI:
dbt run --vars '{"target_date": "2024-01-15", "full_refresh": false}'
В модели:
{% set target_date = var('target_date', none) %}
{{ config(
materialized='incremental',
pre_hook="SET spark.sql.shuffle.partitions = {{ var('shuffle_partitions', 200) }}"
) }}
SELECT *
FROM {{ ref('raw_events') }}
{% if target_date %}
WHERE event_date = '{{ target_date }}'
{% endif %}
Это позволяет параметризовать запуски без изменения кода.
Архитектура кастомных макросов для Spark¶
Рассмотрим реальную библиотеку макросов для Spark-проектов. Хорошо организованные макросы повышают переиспользуемость и упрощают поддержку.
Диаграмма взаимодействия компонентов¶
Jinja compile - это центральная фаза: все конфиги, макросы и глобальные defaults сходятся здесь в конкретный DDL. Spark получает готовый SQL - никакой Jinja на стороне Spark.
Библиотека утилитарных макросов¶
{# macros/spark_utils.sql #}
{# Рассчитать оптимальное число shuffle partitions #}
{% macro optimal_shuffle_partitions(size_gb, target_mb=128) %}
{%- set raw = (size_gb * 1024 / target_mb) | int -%}
{%- set bounded = [50, [raw, 2000] | min] | max -%}
{{ bounded }}
{% endmacro %}
{# Получить location_root в зависимости от target #}
{% macro get_location(layer='gold') %}
{%- if target.name == 'prod' -%}
s3a://prod-lake/{{ layer }}
{%- elif target.name == 'staging' -%}
s3a://staging-lake/{{ layer }}
{%- else -%}
s3a://dev-lake/{{ target.schema }}/{{ layer }}
{%- endif -%}
{% endmacro %}
{# Добавить аудит-колонки к SELECT #}
{% macro add_audit_columns() %}
current_timestamp() AS _loaded_at,
'{{ invocation_id }}' AS _dbt_run_id,
'{{ this.name }}' AS _model_name
{% endmacro %}
{# Создать конфигурацию для Gold-витрины с типичными настройками #}
{% macro gold_table_config(
partition_field,
cluster_field=none,
num_buckets=128,
size_gb=10
) %}
{{
config(
materialized='incremental',
file_format='parquet',
incremental_strategy='insert_overwrite',
partition_by={'field': partition_field, 'data_type': 'date'},
clustered_by=[cluster_field] if cluster_field else none,
buckets=num_buckets if cluster_field else none,
location_root=get_location('gold'),
pre_hook='SET spark.sql.shuffle.partitions = ' ~ optimal_shuffle_partitions(size_gb)
)
}}
{% endmacro %}
Использование gold_table_config в модели:
-- models/marts/fct_transactions.sql
{{ gold_table_config(
partition_field='event_date',
cluster_field='user_id',
num_buckets=256,
size_gb=100
) }}
SELECT
transaction_id,
user_id,
amount,
currency,
{{ add_audit_columns() }},
event_date
FROM {{ ref('int_transactions_enriched') }}
После dbt compile:
-- target/compiled/my_project/models/marts/fct_transactions.sql
CREATE TABLE my_schema.fct_transactions
USING PARQUET
PARTITIONED BY (event_date)
CLUSTERED BY (user_id) INTO 256 BUCKETS
LOCATION 's3a://dev-lake/dev_ivan/gold/fct_transactions'
AS
SELECT
transaction_id,
user_id,
amount,
currency,
current_timestamp() AS _loaded_at,
'abc123-def456' AS _dbt_run_id,
'fct_transactions' AS _model_name,
event_date
FROM my_schema.int_transactions_enriched
Один вызов gold_table_config() заменяет 10+ строк конфигурации и гарантирует единообразие всех Gold-витрин.
Макрос для инкрементального merge с deduplication¶
Частая задача - инкрементальный merge с защитой от дублей:
{# macros/incremental_merge.sql #}
{% macro incremental_merge_config(
unique_key,
partition_field,
cluster_field=none,
num_buckets=128
) %}
{{
config(
materialized='incremental',
file_format='delta',
incremental_strategy='merge',
unique_key=unique_key,
partition_by={'field': partition_field, 'data_type': 'date'},
clustered_by=[cluster_field] if cluster_field else none,
buckets=num_buckets if cluster_field else none,
merge_update_columns=['amount', 'status', 'updated_at', '_loaded_at']
)
}}
{% endmacro %}
Макрос для CDC-пайплайнов¶
CDC (Change Data Capture) - особый случай, где нужно корректно обрабатывать INSERT, UPDATE и DELETE события. Напишем макрос, который инкапсулирует логику CDC-merge:
{# macros/cdc_merge.sql #}
{% macro cdc_incremental(
unique_key,
partition_field,
operation_col='operation',
delete_value='DELETE'
) %}
{{
config(
materialized='incremental',
file_format='delta',
incremental_strategy='merge',
unique_key=unique_key
)
}}
{% if is_incremental() %}
WITH source_data AS (
{{ caller() }}
),
{# Дедупликация: берём последнее событие по каждому ключу #}
deduplicated AS (
SELECT *
FROM (
SELECT
*,
ROW_NUMBER() OVER (
PARTITION BY {{ unique_key }}
ORDER BY updated_at DESC
) AS _rn
FROM source_data
)
WHERE _rn = 1
)
SELECT * FROM deduplicated
{% else %}
{# Первый запуск: только не-удалённые записи #}
WITH source_data AS (
{{ caller() }}
)
SELECT * FROM source_data
WHERE {{ operation_col }} != '{{ delete_value }}'
{% endif %}
{% endmacro %}
Использование:
-- models/marts/dim_customers.sql
{% call cdc_incremental(
unique_key='customer_id',
partition_field='updated_date'
) %}
SELECT
customer_id,
name,
email,
operation,
updated_at,
date(updated_at) AS updated_date
FROM {{ ref('stg_customers_cdc') }}
{% if is_incremental() %}
WHERE updated_at > (SELECT MAX(updated_at) FROM {{ this }})
{% endif %}
{% endcall %}
Макрос call вставляет содержимое блока вместо {{ caller() }} внутри макроса. Это мощный паттерн для оборачивания произвольного SQL в стандартную логику.
Практика: оптимизация витрины fct_transactions¶
Применим всё изученное на реальном примере. Исходная ситуация: витрина fct_transactions обрабатывает события за 4 года, ~100 GB данных. Текущая модель работает медленно: полный скан при каждом запросе, shuffle при join с dim_users.
Шаг 1: Анализ текущего состояния¶
# Скомпилировать текущую модель и посмотреть DDL
dbt compile --select fct_transactions
cat target/compiled/my_project/models/marts/fct_transactions.sql
Текущий DDL (неоптимальный):
CREATE TABLE my_schema.fct_transactions
USING PARQUET
AS SELECT
t.transaction_id,
t.user_id,
t.amount,
t.currency,
u.country_code,
t.event_at,
date(t.event_at) AS event_date
FROM my_schema.int_transactions t
JOIN my_schema.dim_users u ON t.user_id = u.user_id
Проблемы:
- Нет
PARTITIONED BY→ полный scan 146 TB при запросе за один день - Нет
CLUSTERED BY→ shuffle всех данных при join сdim_users - Нет
LOCATION→ managed table в дефолтном warehouse - Нет
pre_hook→ дефолтные 200 shuffle partitions для 100 GB данных
Шаг 2: Создание оптимизированной модели¶
-- models/marts/fct_transactions.sql
{{ config(
materialized='incremental',
file_format='parquet',
incremental_strategy='insert_overwrite',
-- Партиционирование по дате для pruning
partition_by={
"field": "event_date",
"data_type": "date"
},
-- Bucketing по user_id для Shuffle-Free Join с dim_users
clustered_by=['user_id'],
buckets=256,
-- External table в нужном слое хранилища
location_root=env_var('DBT_LOCATION_ROOT', 's3a://dev-lake/gold'),
-- Оптимальные shuffle partitions для ~100 GB
pre_hook=[
"SET spark.sql.shuffle.partitions = 800",
"SET spark.sql.adaptive.enabled = true",
"SET spark.sql.adaptive.skewJoin.enabled = true"
],
-- После загрузки: обновить статистику для Spark CBO
post_hook=[
"ANALYZE TABLE {{ this }} COMPUTE STATISTICS FOR COLUMNS user_id, event_date, amount"
]
) }}
{% if is_incremental() %}
{% set watermark_query %}
SELECT MAX(event_date) FROM {{ this }}
{% endset %}
{% set results = run_query(watermark_query) %}
{% set max_date = results.columns[0].values()[0] %}
{% endif %}
SELECT
t.transaction_id,
t.user_id,
t.amount,
t.currency,
u.country_code,
u.user_segment,
t.event_at,
date(t.event_at) AS event_date,
-- Аудит-колонки
current_timestamp() AS _loaded_at,
'{{ invocation_id }}' AS _dbt_run_id
FROM {{ ref('int_transactions_enriched') }} t
JOIN {{ ref('dim_users') }} u
ON t.user_id = u.user_id
{% if is_incremental() %}
WHERE date(t.event_at) > '{{ max_date }}'
{% endif %}
Шаг 3: dim_users тоже должна быть забакетирована¶
Shuffle-Free Join требует, чтобы обе стороны join были забакетированы совместимо:
-- models/marts/dim_users.sql
{{ config(
materialized='table',
file_format='parquet',
-- Те же 256 buckets по user_id, что и в fct_transactions
clustered_by=['user_id'],
buckets=256,
location_root=env_var('DBT_LOCATION_ROOT', 's3a://dev-lake/gold')
) }}
SELECT
user_id,
name,
email,
country_code,
user_segment,
created_at
FROM {{ ref('stg_users') }}
Теперь обе таблицы имеют CLUSTERED BY (user_id) INTO 256 BUCKETS. Spark выполнит BucketedMergeJoin без shuffle.
Шаг 4: Запуск и проверка¶
# Собрать обе модели (сначала dim_users, потом fct_transactions по DAG)
dbt run --select dim_users fct_transactions
# Проверить скомпилированный DDL
cat target/compiled/my_project/models/marts/fct_transactions.sql
# Проверить партиции
dbt run-operation run_query --args '{query: "SHOW PARTITIONS my_schema.fct_transactions LIMIT 5"}'
# Запустить тесты
dbt test --select fct_transactions
Проверка Shuffle-Free Join в Spark UI:
- Открыть Spark UI → вкладка SQL
- Найти последний job fct_transactions
- В DAG плане должно быть
BucketedMergeJoinвместоExchange + SortMergeJoin
Если план показывает обычный SortMergeJoin с Exchange - значит bucketing не применился. Причины:
- Число buckets не совпадает между таблицами
spark.sql.sources.bucketing.enabled= false (по умолчанию true)- Используется
filter_pushdown, который обходит bucket-join (редко)
Шаг 5: Измерение результата¶
-- До оптимизации (в Spark UI):
-- Scan: 146 TB, Shuffle: 146 TB + 10 GB (dim_users), Duration: 2h 15min
-- После оптимизации:
-- Scan: 100 GB (только вчерашняя партиция), Shuffle: 0 bytes, Duration: 8min
-- Speedup: ~17x
Отладка: когда что-то идёт не так¶
dbt compile - первая линия диагностики¶
Всегда начинайте с dbt compile. Это покажет финальный SQL до отправки на Spark:
dbt compile --select fct_transactions
cat target/compiled/my_project/models/marts/fct_transactions.sql
Если DDL выглядит не так, как ожидалось (нет PARTITIONED BY, нет LOCATION) - проблема в Jinja, не в Spark.
Частые ошибки и их причины¶
Ошибка 1: partition_by не попал в DDL
# Симптом: нет PARTITIONED BY в target/compiled/...
# Причина: опечатка в ключе словаря
{{ config(
partition_by={"field": "event_date", "datatype": "date"} # НЕВЕРНО: datatype
partition_by={"field": "event_date", "data_type": "date"} # ВЕРНО
) }}
Ошибка 2: location создаётся не там
# Симптом: таблица создаётся в /user/hive/warehouse/
# Причина: env_var() вернул пустую строку или None
# Диагностика:
dbt run-operation run_query --args '{query: "SET spark.sql.warehouse.dir"}'
# Проверить переменную окружения:
echo $DBT_LOCATION_ROOT
Ошибка 3: BucketedMergeJoin не применяется
-- Диагностика в Spark:
EXPLAIN SELECT t.*, u.country_code
FROM fct_transactions t
JOIN dim_users u ON t.user_id = u.user_id;
-- В плане ищем:
-- ХОРОШО: BucketedMergeJoin
-- ПЛОХО: Exchange + SortMergeJoin (значит bucketing не применился)
-- Проверить metadata:
DESCRIBE EXTENDED my_schema.fct_transactions;
-- Должно быть: Num Buckets: 256, Bucket Columns: [user_id]
Ошибка 4: submit_timeout не помогает
# submit_timeout контролирует только query submit, не total execution
# Для долгих запросов используйте:
# method: http → submit_timeout: 600 (достаточно для cold start)
# method: thrift → нет аналога, используйте connect_timeout + connect_retries
Макрос для проверки конфигурации таблицы¶
{# macros/describe_table.sql #}
{% macro describe_table(relation) %}
{% set sql %}
DESCRIBE EXTENDED {{ relation }}
{% endset %}
{% do log(run_query(sql).print_table(), info=True) %}
{% endmacro %}
# Использование:
dbt run-operation describe_table --args '{relation: "my_schema.fct_transactions"}'
Выводит полную информацию о таблице: тип, партиции, buckets, location, serde properties.
Итоговая схема принятия решений¶
Схема помогает систематически выбирать нужные параметры конфигурации. Начинаем с размера: таблицы до 10 GB редко требуют специальной оптимизации - Spark справится с дефолтными настройками. Для больших таблиц последовательно задаём себе вопросы: нужен ли pruning по дате, нужен ли Shuffle-Free Join, нужно ли внешнее хранилище, есть ли проблема с таймаутом.
Домашнее задание¶
Задача 1: Диагностика существующей таблицы¶
Возьмите существующую dbt-модель в вашем проекте (или создайте простую тестовую).
- Запустите
dbt compile --select <model_name>и изучите DDL. - Откройте Spark UI после
dbt runи найдите план выполнения. - Ответьте на вопросы:
- Есть ли
PARTITIONED BYв DDL? - Есть ли
CLUSTERED BYв DDL? - Есть ли
Exchange(shuffle) в плане выполнения? - Сколько байт было прочитано со storage?
Задача 2: Добавить partition_by¶
Добавьте partition_by к вашей модели (или fct_transactions из примера).
- Выберите правильный partition key для ваших данных.
- Запустите
dbt run --full-refresh. - Выполните
SHOW PARTITIONS <table>и убедитесь, что партиции создались. - Выполните запрос с фильтром по партиционному полю. Сравните размер прочитанных данных до и после (Spark UI → SQL tab → Scan metrics).
Задача 3: Добавить bucketing и проверить Shuffle-Free Join¶
- Найдите join в вашей модели.
- Добавьте
clustered_byиbucketsк обеим участвующим таблицам с одинаковым ключом и числом buckets. - Запустите
dbt run --full-refreshдля обеих таблиц. - Выполните join в Spark UI. Убедитесь, что план показывает
BucketedMergeJoinвместоSortMergeJoin + Exchange.
Задача 4: Создать библиотеку макросов¶
Создайте файл macros/gold_config.sql с макросом, который:
- Принимает
partition_field,cluster_field,size_gbкак параметры. - Рассчитывает
bucketsкакsize_gb * 5, но не меньше 64 и не больше 2048. - Устанавливает
shuffle_partitionsкакsize_gb * 8черезpre_hook. - Берёт
location_rootизenv_var('DBT_LOCATION_ROOT', 's3a://dev-lake/gold').
Примените макрос к минимум двум моделям и проверьте генерируемый DDL через dbt compile.
Задача 5: Изоляция окружений¶
- Настройте два target в
profiles.yml:devиprod. - Создайте макрос
get_location(layer), который возвращает разные пути для каждого target. - Запустите
dbt compile --target devиdbt compile --target prod. - Убедитесь, что
LOCATIONв DDL отличается между окружениями.
Чеклист перед выпуском модели¶
Прежде чем считать модель готовой к production, пройдите по этому списку:
- [ ] partition_by добавлен для таблиц > 10 GB?
- [ ] data_type в partition_by указан корректно (не timestamp вместо date)?
- [ ] clustered_by добавлен, если есть активный JOIN по фиксированному ключу?
- [ ] buckets - степень двойки и совместима с другой таблицей в JOIN?
- [ ] location_root указывает на правильный слой (bronze/silver/gold)?
- [ ] location_root использует
env_var()для изоляции окружений? - [ ] submit_timeout в profiles.yml достаточен для cold start?
- [ ] server_side_parameters включают AQE (
spark.sql.adaptive.enabled: "true")? - [ ] dbt compile показывает ожидаемый DDL?
- [ ] dbt test проходит после
dbt run? - [ ] Spark UI показывает
BucketedMergeJoin(если добавлен bucketing)? - [ ] SHOW PARTITIONS показывает ожидаемые партиции?