dbt-spark адаптер: методы подключения thrift, http, session и profiles.yml
dbt-spark - адаптер, соединяющий dbt Core и Apache Spark. Разбираем методы thrift, http, session, анатомию profiles.yml, передачу Spark-конфигов и интеграцию с Delta Lake и MinIO.
Архитектурный мост между dbt и Spark¶
Где заканчивается dbt и начинается Spark¶
dbt Core - это инструмент, который работает на машине разработчика или в CI-контейнере. Он читает SQL-файлы моделей, компилирует Jinja-шаблоны, строит граф зависимостей и отправляет готовые SQL-запросы на выполнение. Но сам dbt Core не умеет исполнять SQL - он только генерирует его и управляет порядком выполнения.
Для исполнения SQL-кода dbt нужен адаптер - плагин, который знает, как общаться с конкретной СУБД или вычислительным движком. Адаптер переводит внутренние абстракции dbt в вызовы конкретного API и диалекта SQL.
Для Apache Spark этот адаптер - dbt-spark. Его задача: принять от dbt Core скомпилированный SQL-запрос (например, CREATE TABLE dbt_dev.stg_events AS SELECT ...) и выполнить его через одно из доступных соединений со Spark-кластером.
Жизненный цикл запроса dbt-spark¶
Разберём путь от .sql-файла до выполненного Spark-запроса:
Разработчик
│
├── [Локально] dbt compile
│ ├── Читает models/staging/stg_events.sql
│ ├── Компилирует Jinja: {{ ref(...) }} → реальные имена
│ └── Сохраняет в target/compiled/
│
├── [Локально] dbt run
│ ├── Строит DAG из зависимостей
│ ├── Для каждой модели:
│ │ ├── Читает compiled SQL
│ │ ├── Оборачивает в DDL (CREATE TABLE / MERGE INTO / ...)
│ │ └── Передаёт в адаптер dbt-spark
│ └── Ждёт результата от адаптера
│
└── [Адаптер dbt-spark]
├── Устанавливает соединение (thrift/http/session)
├── Отправляет SQL на Spark-кластер
├── Получает статус выполнения
└── Возвращает результат в dbt Core
Ключевое разделение: компиляция происходит локально (на машине разработчика или в CI-runner'е), выполнение - на Spark-кластере. Это означает, что dbt-процесс не нагружает кластер компиляцией шаблонов - он отправляет уже готовый SQL.
Роль адаптера: что именно он делает¶
Адаптер dbt-spark выполняет несколько видов операций:
DDL-команды: создание, замена, удаление таблиц и views. Разные форматы хранения (Parquet, Delta, Iceberg) требуют разных DDL-конструкций, и адаптер знает, как правильно сформировать CREATE TABLE ... USING delta PARTITIONED BY ....
DML-команды: вставка данных, обновление (MERGE INTO), перезапись партиций (INSERT OVERWRITE). Адаптер генерирует правильный диалект для конкретной версии Spark.
Метаданные: получение списка таблиц, проверка существования объектов, получение схемы таблицы. dbt использует это для проверок перед выполнением модели.
Управление сессией: установка Spark-конфигураций (shuffle.partitions, AQE), управление временем жизни соединения.
Установка и зависимости¶
Адаптер dbt-spark - это Python-пакет с опциональными зависимостями в зависимости от метода подключения:
pip install "dbt-spark[PyHive]"
pip install "dbt-spark[ODBC]"
pip install "dbt-spark[all]"
pip install "dbt-spark"
Опциональные зависимости:
| Метод | Зависимости | Установка |
|---|---|---|
thrift |
PyHive, thrift, thrift-sasl | dbt-spark[PyHive] |
http |
requests, urllib3 | Встроены в dbt-spark |
session |
pyspark | pip install pyspark отдельно |
odbc |
pyodbc | dbt-spark[ODBC] |
pip install "dbt-spark[PyHive]" pyspark==3.5.0
dbt --version
Core:
- installed: 1.7.3
- latest: 1.7.3 - Up to date!
Plugins:
- spark: 1.7.1 - Up to date!
Архитектура подключения¶
Диаграмма показывает полную картину: dbt Core (слева вверху) отправляет SQL через один из четырёх методов подключения. Каждый метод приходит к своей точке входа на стороне кластера, но все они в итоге выполняют SQL через одни и те же Spark Executors, которые читают/пишут данные в объектное хранилище и обновляют метаданные в Hive Metastore или Iceberg Catalog.
Выбор метода подключения определяется топологией сети, требованиями безопасности и типом Spark-окружения - а не тем, "что лучше по качеству SQL". Все четыре метода выполняют одинаковый SQL на одном и том же Spark-движке.
Метод thrift: классическое подключение¶
Что такое Spark Thrift Server¶
Spark Thrift Server (официальное название в документации Spark - Spark SQL Thrift Server) - это компонент Spark, который превращает кластер в SQL-сервис с доступом по протоколу JDBC/ODBC. Он реализует тот же интерфейс, что Apache Hive HiveServer2 (HiveServer2 - это "SQL-фронтенд" Hive, принимающий запросы через порт 10000).
Spark Thrift Server запускается как долгоживущий процесс на узле кластера:
$SPARK_HOME/sbin/start-thriftserver.sh \
--master spark://spark-master:7077 \
--conf spark.sql.hive.thriftServer.singleSession=false \
--conf spark.sql.warehouse.dir=s3a://datalake/warehouse/ \
--hiveconf hive.server2.thrift.port=10000 \
--hiveconf hive.server2.thrift.bind.host=0.0.0.0
После запуска клиенты могут подключаться через:
- Beeline (CLI Hive):
beeline -u jdbc:hive2://spark-thrift:10000/default - JDBC: стандартный JDBC-драйвер Hive
- ODBC: драйвер Simba или Cloudera Hive ODBC
- dbt-spark через метод
thrift - BI-инструменты: Tableau, Superset, Metabase - те, что умеют подключаться к Hive
Концепция HiveServer2 важна для понимания: Spark Thrift Server реализует тот же протокол Apache Thrift (бинарный RPC-фреймворк), что и оригинальный HiveServer2. Поэтому любой клиент, умеющий работать с Hive, работает и со Spark Thrift Server.
Конфигурация thrift в profiles.yml¶
# ~/.dbt/profiles.yml
analytics:
target: dev
outputs:
dev:
type: spark
method: thrift
host: spark-thrift.internal
port: 10000
user: dbt_service
password: "{{ env_var('SPARK_THRIFT_PASSWORD', '') }}"
schema: dbt_dev
connect_timeout: 60
connect_retries: 3
threads: 4
auth: CUSTOM
kerberos_service_name: hive
use_ssl: false
server_side_parameters:
spark.sql.shuffle.partitions: "200"
spark.sql.adaptive.enabled: "true"
Разберём все параметры:
host - имя хоста или IP-адрес Spark Thrift Server. В production это обычно DNS-имя, не IP.
port - порт. Стандарт для HiveServer2 - 10000. Если Thrift Server настроен на другой порт, укажите соответственно.
user и password - учётные данные. Важно: пароль никогда не должен храниться в профиле как plain text. Используйте env_var() - об этом подробно в разделе о безопасности.
schema - в контексте Spark это имя базы данных (namespace). Все таблицы, которые dbt создаёт в этом профиле, будут в этой базе. На dev это dbt_dev, на prod - dbt_prod (или analytics, gold и т.д.).
connect_timeout - сколько секунд ждать установки соединения. Если Spark Thrift Server работает на кластере с autoscaling, первый запрос может ждать, пока кластер "проснётся". 60 секунд - разумный таймаут для кластеров с cold start.
connect_retries - количество повторных попыток соединения при таймауте или временной недоступности. Вместе с connect_timeout это обеспечивает устойчивость пайплайна к кратковременным сбоям сети.
threads - количество параллельных потоков dbt. Каждый поток - это одно соединение со Spark Thrift Server. Если threads=4, dbt может выполнять 4 независимых модели одновременно. Не путайте с числом Spark-задач внутри одного запроса: Spark всегда параллелен внутри одного запроса независимо от этого параметра.
auth - метод аутентификации:
NONE- без аутентификации (только для dev, никогда для prod)CUSTOM- логин/пароль через SASLKERBEROS- Kerberos-аутентификация для enterprise-кластеров
server_side_parameters - словарь Spark-конфигураций, которые dbt устанавливает в начале каждой сессии. Это эквивалент spark.conf.set(...) в PySpark - только передаётся через JDBC сразу при подключении.
Когда использовать thrift¶
Thrift - правильный выбор когда:
- У вас self-hosted Spark кластер (bare metal, Kubernetes, Yarn, Mesos)
- Вам нужна многопользовательская среда (несколько аналитиков работают с одним кластером)
- Вы интегрируете dbt с BI-инструментами (Superset, Tableau), которые тоже подключаются через HiveServer2
- Нужна совместимость с Hive - существующие инструменты, настроенные на HiveServer2, продолжают работать
Thrift - не лучший выбор когда:
- Кластер часто останавливается (autoscaling, spot-instances) - каждый cold start потребует ожидания
- Много коротких независимых dbt-запусков - накладные расходы на установку соединения заметны
- Нужна облачная управляемость - у HTTP-метода лучше интеграция с IAM, токенами и API-шлюзами
Метод http: облачный стандарт¶
Что такое HTTP-подключение к Spark¶
Метод http используется для подключения к Spark через HTTP/HTTPS REST-эндпоинт. Это характерно для облачных управляемых Spark-сервисов:
- Databricks SQL Warehouse - управляемый SQL-сервис Databricks с REST API
- AWS EMR Studio - Jupyter-based сервис с HTTP-доступом к EMR
- Google Dataproc - managed Spark на GCP с HTTP-gateway
- Любой HTTP-прокси поверх Spark с REST API
HTTP-метод не устанавливает постоянное TCP-соединение (как thrift), а отправляет каждый запрос как отдельный HTTP POST, получает статус выполнения через polling и забирает результат.
Конфигурация http в profiles.yml¶
analytics:
target: databricks
outputs:
databricks:
type: spark
method: http
schema: dbt_prod
host: dbc-xxxxxxxx-xxxx.cloud.databricks.com
endpoint: /sql/1.0/warehouses/xxxxxxxxxxxxxxxx
token: "{{ env_var('DATABRICKS_TOKEN') }}"
port: 443
use_ssl: true
connect_timeout: 120
connect_retries: 5
threads: 8
server_side_parameters:
spark.sql.shuffle.partitions: "1000"
spark.sql.adaptive.enabled: "true"
spark.sql.adaptive.advisoryPartitionSizeInBytes: "134217728"
Ключевые параметры HTTP-метода:
host - FQDN облачного сервиса. Для Databricks это <workspace-id>.cloud.databricks.com.
endpoint - путь к SQL-warehouse'у или конкретному кластеру. Для Databricks SQL Warehouse: /sql/1.0/warehouses/<warehouse-id>. Этот параметр отличает HTTP-метод от других: он определяет конкретный compute-ресурс, к которому направляются запросы.
token - Bearer-токен для аутентификации. В Databricks это Personal Access Token (PAT) или Service Principal токен. Никогда не хардкодьте токен в файл - только через env_var().
port: 443 - стандартный HTTPS-порт.
use_ssl: true - шифрование трафика. Для публичных облачных сервисов всегда true.
threads: 8 - HTTP-метод хорошо масштабируется по потокам, потому что каждый запрос независимый. Databricks SQL Warehouse может параллельно обрабатывать много запросов из разных потоков dbt.
Схема аутентификации HTTP¶
dbt-spark Databricks/Cloud Service
│ │
├── HTTP POST /sql/execute ─→ │ (Bearer: <token>)
│ {sql: "CREATE TABLE ..."} │
│ │
│ ← HTTP 202 Accepted ────── │ (statement_id: "xxx")
│ │
├── HTTP GET /statements/xxx → │ (polling статус)
│ ← HTTP 200 {state: "RUNNING"}
│
├── HTTP GET /statements/xxx → │
│ ← HTTP 200 {state: "SUCCEEDED", result: {...}}
│
└── Готово
Polling-подход означает, что при долгих запросах (например, полный пересчёт большой таблицы) dbt периодически опрашивает статус, не блокируя потребление ресурсов на клиентской стороне.
Когда использовать http¶
HTTP - правильный выбор когда:
- Вы работаете с управляемым облачным Spark (Databricks, EMR, Dataproc)
- Нужна аутентификация через токены и интеграция с корпоративным IdP (SSO, SAML)
- Кластер доступен только через HTTPS (нет прямого TCP-доступа к порту 10000)
- Нужна горизонтальная масштабируемость: облачные SQL warehouse'ы автоматически масштабируются под нагрузку
Метод session: встроенный Spark¶
Идея session-метода¶
Метод session - принципиально другой подход. Вместо подключения к удалённому Spark-кластеру, dbt запускает SparkSession прямо внутри своего Python-процесса. Никакого сетевого соединения нет - Spark и dbt работают в одном процессе на одной машине.
Это возможно, потому что PySpark позволяет запустить полноценный Spark Driver локально: данные обрабатываются на одной машине (или в одном контейнере), без дополнительных Executor-узлов. Для небольших датасетов и тестов этого достаточно.
Конфигурация session в profiles.yml¶
analytics:
target: local
outputs:
local:
type: spark
method: session
schema: dbt_local
threads: 2
server_side_parameters:
spark.executor.memory: "2g"
spark.driver.memory: "4g"
spark.sql.shuffle.partitions: "4"
spark.sql.adaptive.enabled: "true"
spark.master: "local[*]"
spark.sql.extensions: >
io.delta.sql.DeltaSparkSessionExtension
spark.sql.catalog.spark_catalog: >
org.apache.spark.sql.delta.catalog.DeltaCatalog
spark.hadoop.fs.s3a.endpoint: "http://localhost:9000"
spark.hadoop.fs.s3a.access.key: "{{ env_var('MINIO_ACCESS_KEY') }}"
spark.hadoop.fs.s3a.secret.key: "{{ env_var('MINIO_SECRET_KEY') }}"
spark.hadoop.fs.s3a.path.style.access: "true"
Разберём специфику session-конфигурации:
spark.master: "local[*]" - запуск Spark в локальном режиме с использованием всех доступных ядер CPU. local[4] - 4 ядра, local[*] - все ядра.
spark.driver.memory и spark.executor.memory - в локальном режиме driver и executor - это один процесс. Оба параметра влияют на доступный heap.
spark.sql.shuffle.partitions: "4" - для локальных тестов 4 партиции достаточно. Дефолт 200 создаёт 200 пустых задач на крошечных датасетах.
Delta Lake расширения - подключаются через spark.sql.extensions и spark.sql.catalog.spark_catalog. В session-режиме это критично: без этих конфигураций команды CREATE TABLE ... USING delta не будут работать.
MinIO интеграция - параметры spark.hadoop.fs.s3a.* позволяют session-режиму читать/писать данные в локальный MinIO. Это делает session-режим полноценной средой для разработки с реальными данными, но без удалённого кластера.
Как dbt запускает SparkSession¶
При методе session адаптер dbt-spark создаёт SparkSession при инициализации:
from pyspark.sql import SparkSession
spark = SparkSession.builder \
.master("local[*]") \
.config("spark.driver.memory", "4g") \
.config("spark.sql.shuffle.partitions", "4") \
.config("spark.sql.extensions",
"io.delta.sql.DeltaSparkSessionExtension") \
.config("spark.sql.catalog.spark_catalog",
"org.apache.spark.sql.delta.catalog.DeltaCatalog") \
.getOrCreate()
После создания сессии все SQL-запросы от dbt исполняются через spark.sql(query). SparkSession живёт всё время выполнения dbt run и завершается вместе с процессом.
session vs thrift: принципиальная разница¶
| Аспект | session | thrift |
|---|---|---|
| Где работает Spark | На машине разработчика / CI-runner | На удалённом кластере |
| Сетевое соединение | Нет | Да (TCP порт 10000) |
| Масштабируемость | Ограничена одной машиной | Полноценный кластер |
| Latency | Нет накладных расходов на соединение | Есть (мс–сек) |
| Изоляция | Полная (свой процесс) | Shared (общий Thrift Server) |
| Требования | PySpark установлен локально | Thrift Server запущен на кластере |
| Типичное использование | Dev / CI / Unit tests | Production / Shared analytics |
Когда использовать session¶
Session - правильный выбор когда:
- Локальная разработка: хотите быстро проверить модель без подключения к кластеру
- CI/CD тесты: GitHub Actions / GitLab CI - каждый runner поднимает свой Spark
- Unit-тестирование: тесты с фиксированными мини-датасетами
- Open-source lakehouse на одной машине: разработчик работает с MinIO + Delta Lake локально
Session - не подходит когда:
- Данных больше, чем влезает на одну машину (терабайты)
- Нужна конкурентная работа нескольких пользователей
- Production-нагрузка - одна машина не может заменить кластер
Анатомия profiles.yml¶
Где хранится и как загружается profiles.yml¶
По умолчанию dbt ищет profiles.yml в ~/.dbt/profiles.yml (домашняя директория пользователя). Это намеренно вынесено из репозитория: профили содержат учётные данные и не должны попадать в Git.
Порядок поиска:
- Переменная окружения
DBT_PROFILES_DIR- если задана, dbt ищет в этой директории - Аргумент
--profiles-dirпри запуске команды - Файл
profiles.ymlв текущей директории проекта ~/.dbt/profiles.yml- дефолт
export DBT_PROFILES_DIR=/opt/dbt/profiles/
dbt run --profiles-dir /opt/dbt/profiles/
cat profiles.yml
dbt run
В Docker/Kubernetes profiles.yml монтируется как Secret или ConfigMap:
# Kubernetes Secret
apiVersion: v1
kind: Secret
metadata:
name: dbt-profiles
data:
profiles.yml: <base64-encoded-content>
---
# Pod Volume Mount
volumes:
- name: dbt-profiles
secret:
secretName: dbt-profiles
volumeMounts:
- name: dbt-profiles
mountPath: /home/dbt/.dbt
readOnly: true
Полная структура profiles.yml¶
# ~/.dbt/profiles.yml
# Имя профиля - должно совпадать с "profile" в dbt_project.yml
analytics_lakehouse:
# Активный target по умолчанию (можно переопределить --target)
target: dev
# Все доступные окружения
outputs:
# Dev-окружение: локальный session
dev:
type: spark
method: session
schema: dbt_dev_{{ env_var('USER', 'unknown') }}
threads: 2
server_side_parameters:
spark.master: "local[4]"
spark.driver.memory: "4g"
spark.sql.shuffle.partitions: "8"
spark.sql.adaptive.enabled: "true"
spark.sql.extensions: "io.delta.sql.DeltaSparkSessionExtension"
spark.sql.catalog.spark_catalog: >
org.apache.spark.sql.delta.catalog.DeltaCatalog
spark.hadoop.fs.s3a.endpoint: "http://localhost:9000"
spark.hadoop.fs.s3a.access.key: "{{ env_var('MINIO_ACCESS_KEY') }}"
spark.hadoop.fs.s3a.secret.key: "{{ env_var('MINIO_SECRET_KEY') }}"
spark.hadoop.fs.s3a.path.style.access: "true"
# CI-окружение: session с тестовыми данными
ci:
type: spark
method: session
schema: "dbt_ci_{{ env_var('CI_PIPELINE_ID', 'local') }}"
threads: 4
server_side_parameters:
spark.master: "local[*]"
spark.driver.memory: "8g"
spark.sql.shuffle.partitions: "4"
spark.sql.adaptive.enabled: "true"
spark.sql.extensions: "io.delta.sql.DeltaSparkSessionExtension"
spark.sql.catalog.spark_catalog: >
org.apache.spark.sql.delta.catalog.DeltaCatalog
# Staging-окружение: подключение к тестовому кластеру
staging:
type: spark
method: thrift
host: "{{ env_var('SPARK_THRIFT_HOST_STAGING') }}"
port: 10000
user: "{{ env_var('SPARK_USER') }}"
password: "{{ env_var('SPARK_PASSWORD_STAGING') }}"
schema: dbt_staging
auth: CUSTOM
connect_timeout: 60
connect_retries: 3
threads: 4
server_side_parameters:
spark.sql.shuffle.partitions: "400"
spark.sql.adaptive.enabled: "true"
spark.sql.adaptive.advisoryPartitionSizeInBytes: "134217728"
# Production: основной кластер
prod:
type: spark
method: thrift
host: "{{ env_var('SPARK_THRIFT_HOST_PROD') }}"
port: "{{ env_var('SPARK_THRIFT_PORT', '10000') | int }}"
user: "{{ env_var('SPARK_USER_PROD') }}"
password: "{{ env_var('SPARK_PASSWORD_PROD') }}"
schema: analytics
auth: CUSTOM
use_ssl: true
connect_timeout: 120
connect_retries: 5
threads: 8
server_side_parameters:
spark.sql.shuffle.partitions: "2000"
spark.sql.adaptive.enabled: "true"
spark.sql.adaptive.advisoryPartitionSizeInBytes: "134217728"
spark.sql.adaptive.coalescePartitions.enabled: "true"
spark.sql.sources.partitionOverwriteMode: "dynamic"
spark.sql.parquet.compression.codec: "snappy"
Переключение между окружениями¶
dbt run
dbt run --target dev
dbt run --target staging
dbt run --target prod
export DBT_TARGET=staging
dbt run
schema: база данных или namespace¶
В SQL-мире разных систем слово "schema" означает разное. В PostgreSQL schema - это пространство имён внутри базы данных (аналог namespace). В dbt в общем смысле schema - целевая схема для объектов.
В Apache Spark слово "schema" в profiles.yml соответствует database/namespace (то, что в SQL называется USE database_name). Таблица stg_events при schema: dbt_dev создаётся как dbt_dev.stg_events.
Это важно понимать при работе с трёхчастным именованием в Spark 3.x:
-- Трёхчастное имя: catalog.database.table
SELECT * FROM spark_catalog.dbt_dev.stg_events;
-- Двухчастное (без каталога): database.table
SELECT * FROM dbt_dev.stg_events;
threads: параллелизм на уровне dbt¶
Параметр threads определяет, сколько моделей dbt может выполнять одновременно. Каждый поток - одно соединение с кластером и один активный SQL-запрос.
threads: 4 при DAG:
A → C
B → C → D
Волна 1: A, B (параллельно)
Волна 2: C (когда A и B завершены)
Волна 3: D
Правило выбора threads:
session: 1–2 потока (единственный Spark Driver делит ресурсы между ними)thriftна shared-кластере: 4–8 потоков (не перегружать Thrift Server)thriftна dedicated кластере: 8–16 потоковhttpна управляемом облачном сервисе: 8–16 (облачный warehouse масштабируется)
Увеличение threads не всегда ускоряет работу. Если DAG линейный (каждая модель зависит от предыдущей), параллельных задач нет и threads > 1 бесполезны. Threads помогают только при наличии независимых веток в DAG.
Безопасность: секреты в profiles.yml¶
Никогда не хардкодить учётные данные¶
Главное правило работы с profiles.yml: никаких паролей, токенов и ключей в plain text в файле.
# ПЛОХО - хардкод в файле
prod:
password: "super_secret_password_123"
token: "dapi1234567890abcdef"
spark.hadoop.fs.s3a.secret.key: "my-s3-secret"
Если такой файл попадёт в Git (случайно, или в публичный репозиторий), учётные данные скомпрометированы. Даже если потом удалить, данные остаются в истории Git.
env_var() - правильный способ¶
# ПРАВИЛЬНО - все секреты через переменные окружения
prod:
password: "{{ env_var('SPARK_PASSWORD_PROD') }}"
token: "{{ env_var('DATABRICKS_TOKEN') }}"
server_side_parameters:
spark.hadoop.fs.s3a.secret.key: "{{ env_var('S3A_SECRET_KEY') }}"
Синтаксис env_var():
# Без дефолтного значения - ошибка если переменная не задана
"{{ env_var('REQUIRED_VAR') }}"
# С дефолтным значением - использует дефолт если переменная не задана
"{{ env_var('OPTIONAL_VAR', 'default_value') }}"
Переменные окружения устанавливаются:
export SPARK_PASSWORD_PROD="actual_password_here"
export DATABRICKS_TOKEN="dapi1234567890abcdef"
export S3A_SECRET_KEY="my-s3-secret"
dbt run --target prod
В CI/CD системах (GitHub Actions, GitLab CI, Airflow):
# GitHub Actions
env:
SPARK_PASSWORD_PROD: ${{ secrets.SPARK_PASSWORD_PROD }}
DATABRICKS_TOKEN: ${{ secrets.DATABRICKS_TOKEN }}
S3A_SECRET_KEY: ${{ secrets.S3A_SECRET_KEY }}
# GitLab CI
variables:
SPARK_PASSWORD_PROD: $SPARK_PASSWORD_PROD
DATABRICKS_TOKEN: $DATABRICKS_TOKEN
Что хранить в Git vs что в секретах¶
| В Git (в profiles.yml) | В секретах (env_var) |
|---|---|
| host, port | password, token |
| method, auth | access_key, secret_key |
| schema, threads | kerberos_keytab path |
| connect_timeout | any API keys |
| server_side_parameters | certificates, SSL keys |
.gitignore для profiles.yml¶
Если profiles.yml вдруг оказывается в директории проекта (не в ~/.dbt/):
# .gitignore
profiles.yml
.dbt/
target/
logs/
В репозитории dbt-проекта можно хранить profiles.yml.example с шаблоном без реальных значений:
# profiles.yml.example - БЕЗОПАСНО хранить в Git
analytics_lakehouse:
target: dev
outputs:
dev:
type: spark
method: session
schema: dbt_dev
# ... non-secret params ...
prod:
type: spark
method: thrift
host: "{{ env_var('SPARK_THRIFT_HOST_PROD') }}"
password: "{{ env_var('SPARK_PASSWORD_PROD') }}"
# ... non-secret params ...
Spark-конфигурации в profiles.yml¶
server_side_parameters: передача Spark-конфигов¶
server_side_parameters - словарь конфигураций Spark, которые dbt устанавливает при подключении. Для thrift-метода это выполняется через SET spark.conf.key=value в начале JDBC-сессии. Для session-метода - через SparkConf при создании SparkSession.
server_side_parameters:
# Shuffle оптимизация
spark.sql.shuffle.partitions: "800"
spark.sql.adaptive.enabled: "true"
spark.sql.adaptive.advisoryPartitionSizeInBytes: "134217728"
spark.sql.adaptive.coalescePartitions.enabled: "true"
spark.sql.adaptive.skewJoin.enabled: "true"
# Запись партиций
spark.sql.sources.partitionOverwriteMode: "dynamic"
# Delta Lake
spark.sql.extensions: "io.delta.sql.DeltaSparkSessionExtension"
spark.sql.catalog.spark_catalog: >
org.apache.spark.sql.delta.catalog.DeltaCatalog
spark.databricks.delta.retentionDurationCheck.enabled: "false"
# Iceberg (альтернатива Delta)
# spark.sql.extensions: >
# org.apache.iceberg.spark.extensions.IcebergSparkSessionExtensions
# spark.sql.catalog.iceberg: org.apache.iceberg.spark.SparkCatalog
# spark.sql.catalog.iceberg.type: hadoop
# spark.sql.catalog.iceberg.warehouse: s3a://datalake/iceberg/
# Object Storage (S3/MinIO)
spark.hadoop.fs.s3a.endpoint: "http://minio:9000"
spark.hadoop.fs.s3a.access.key: "{{ env_var('MINIO_ACCESS_KEY') }}"
spark.hadoop.fs.s3a.secret.key: "{{ env_var('MINIO_SECRET_KEY') }}"
spark.hadoop.fs.s3a.path.style.access: "true"
spark.hadoop.fs.s3a.impl: "org.apache.hadoop.fs.s3a.S3AFileSystem"
# Память (только для session-метода)
# spark.driver.memory: "8g"
# spark.executor.memory: "4g"
Разные конфигурации для dev и prod¶
outputs:
dev:
type: spark
method: session
server_side_parameters:
spark.sql.shuffle.partitions: "8"
spark.master: "local[4]"
prod:
type: spark
method: thrift
server_side_parameters:
spark.sql.shuffle.partitions: "2000"
spark.sql.adaptive.enabled: "true"
В dev минимальные настройки для быстрой итерации. В prod - тщательно подобранные конфигурации для работы с терабайтными датасетами. Такое разделение гарантирует, что разработчики случайно не запустят тяжёлые конфигурации на локальной машине.
dbt debug: диагностика подключения¶
Что проверяет dbt debug¶
Команда dbt debug выполняет полную диагностику окружения dbt без выполнения моделей:
dbt debug
15:42:33 Running with dbt=1.7.3
15:42:33 dbt version: 1.7.3
15:42:33 python version: 3.11.6
15:42:33 python path: /usr/bin/python3
15:42:33 os info: Linux-6.1.0-1022-aws-x86_64
15:42:33 Using profiles.yml file at /home/ubuntu/.dbt/profiles.yml
15:42:33 Using dbt_project.yml file at /home/ubuntu/analytics_lakehouse/dbt_project.yml
15:42:33 Configuration:
15:42:33 profiles.yml file [OK found and valid]
15:42:33 dbt_project.yml file [OK found and valid]
15:42:33 Required dependencies:
15:42:33 - git [OK found]
15:42:33 Connection:
15:42:35 host: spark-thrift.internal
15:42:35 port: 10000
15:42:35 database: None
15:42:35 schema: dbt_dev
15:42:35 auth: CUSTOM
15:42:35 Connection test: [OK connection ok]
15:42:35 All checks passed!
Последняя строка Connection test: [OK connection ok] означает, что dbt успешно подключился к Spark Thrift Server и выполнил тестовый запрос.
Типичные ошибки и их диагностика¶
Ошибка: профиль не найден
Runtime Error
Could not find profile named 'analytics_lakehouse'
Причина: в dbt_project.yml указан profile: analytics_lakehouse, но в profiles.yml нет такого профиля. Проверьте написание.
Ошибка: переменная окружения не задана
Compilation Error in profiles.yml
The following variables were undefined:
SPARK_PASSWORD_PROD (referenced in profiles.yml)
Причина: env_var('SPARK_PASSWORD_PROD') без дефолтного значения, а переменная не установлена. Решение: export SPARK_PASSWORD_PROD="..." перед запуском.
Ошибка: соединение отклонено
Connection test: [ERROR]
Could not connect to server: Connection refused (0.0.0.0:10000)
Причина: Spark Thrift Server не запущен или неправильный host/port. Проверьте ps aux | grep ThriftServer на кластере.
Ошибка: аутентификация
Connection test: [ERROR]
TTransportException(type=1, message='Could not connect to ...')
Error: (): Could not connect to xxx (code SASL_ERROR)
Причина: неверные учётные данные или неправильный auth (NONE вместо CUSTOM). Проверьте логи Thrift Server.
dbt debug --log-level debug 2>&1 | grep -i "error\|exception\|connection"
Интеграция с Delta Lake через session¶
Полная настройка Delta Lake в session-режиме¶
# profiles.yml - dev с Delta Lake + MinIO
analytics:
target: dev
outputs:
dev:
type: spark
method: session
schema: dbt_dev
threads: 2
server_side_parameters:
spark.master: "local[4]"
spark.driver.memory: "4g"
spark.jars.packages: >
io.delta:delta-core_2.12:2.4.0,
org.apache.hadoop:hadoop-aws:3.3.4,
com.amazonaws:aws-java-sdk-bundle:1.12.262
spark.sql.extensions: >
io.delta.sql.DeltaSparkSessionExtension
spark.sql.catalog.spark_catalog: >
org.apache.spark.sql.delta.catalog.DeltaCatalog
spark.hadoop.fs.s3a.endpoint: "http://localhost:9000"
spark.hadoop.fs.s3a.access.key: "minioadmin"
spark.hadoop.fs.s3a.secret.key: "{{ env_var('MINIO_SECRET') }}"
spark.hadoop.fs.s3a.path.style.access: "true"
spark.hadoop.fs.s3a.impl: >
org.apache.hadoop.fs.s3a.S3AFileSystem
spark.hadoop.fs.s3a.aws.credentials.provider: >
org.apache.hadoop.fs.s3a.SimpleAWSCredentialsProvider
spark.sql.shuffle.partitions: "4"
spark.sql.adaptive.enabled: "true"
spark.jars.packages - автоматическое скачивание JAR-зависимостей из Maven Central при запуске SparkSession. Это удобно для development, но в production лучше предварительно скопировать JAR в /opt/spark/jars/.
Модель с Delta-специфичными настройками¶
-- models/marts/fct_daily_revenue.sql
{{ config(
materialized='incremental',
incremental_strategy='merge',
unique_key=['campaign_id', 'event_date'],
file_format='delta',
location='s3a://datalake/gold/fct_daily_revenue',
partition_by=['event_date'],
tblproperties={
'delta.autoOptimize.optimizeWrite': 'true',
'delta.autoOptimize.autoCompact': 'true',
'delta.logRetentionDuration': 'interval 30 days',
}
) }}
SELECT
campaign_id,
CAST(event_ts AS DATE) AS event_date,
SUM(revenue) AS total_revenue,
COUNT(*) AS event_count,
CURRENT_TIMESTAMP() AS updated_at
FROM {{ ref('stg_events') }}
{% if is_incremental() %}
WHERE CAST(event_ts AS DATE) >= DATE_ADD(CURRENT_DATE(), -3)
{% endif %}
GROUP BY campaign_id, CAST(event_ts AS DATE)
tblproperties - свойства Delta-таблицы. delta.autoOptimize.optimizeWrite автоматически оптимизирует размер файлов при записи (важно для предотвращения small files problem). delta.autoOptimize.autoCompact запускает фоновую компакцию после записи.
Проверка Delta-таблицы после запуска¶
DESCRIBE EXTENDED dbt_dev.fct_daily_revenue;
DESCRIBE HISTORY dbt_dev.fct_daily_revenue;
SELECT * FROM dbt_dev.fct_daily_revenue@v1;
VACUUM dbt_dev.fct_daily_revenue RETAIN 168 HOURS;
Мониторинг: Spark UI и dbt logs¶
Как dbt-модели отображаются в Spark UI¶
Когда dbt выполняет модель, в Spark UI создаётся Job с именем, которое включает название модели. Это позволяет коррелировать медленные dbt-модели с конкретными Spark-job'ами.
Spark UI → Jobs:
Job 0: "CREATE TABLE dbt_dev.stg_events AS SELECT ..."
Job 1: "MERGE INTO dbt_dev.fct_user_activity ..."
Job 2: "CREATE VIEW dbt_dev.stg_users AS SELECT ..."
Для диагностики медленных моделей:
- В dbt логах найдите, какая модель долго выполняется
- В Spark UI найдите соответствующий Job по времени или имени таблицы
- Проверьте метрики Stage'ей: Shuffle Read Size, Spill, Task Duration
- Скорректируйте
server_side_parametersв profiles.yml
dbt logs для диагностики¶
cat logs/dbt.log | grep -E "Running|Completed|Error"
cat logs/dbt.log | grep "fct_user_activity"
cat logs/dbt.log | grep -i "merge into"
cat logs/dbt.log | grep -E "ERROR|WARNING"
Полные SQL-запросы, которые dbt отправляет в Spark, видны в logs/dbt.log при уровне логирования DEBUG:
DBT_LOG_LEVEL=debug dbt run --select fct_user_activity 2>&1 | \
grep "SQL Query\|On model"
Лабораторная практика¶
Сценарий: настройка dev и prod окружений¶
Настроим dbt-проект с двумя окружениями:
- dev: session-режим с локальным MinIO и Delta Lake
- prod: thrift-подключение к кластерному Spark Thrift Server
Шаг 1: Запуск Spark Thrift Server в Docker¶
# docker-compose.yml
version: "3.8"
services:
minio:
image: minio/minio:latest
ports:
- "9000:9000"
- "9001:9001"
environment:
MINIO_ROOT_USER: minioadmin
MINIO_ROOT_PASSWORD: minioadmin123
command: server /data --console-address ":9001"
volumes:
- minio_data:/data
spark-thrift:
image: bitnami/spark:3.5
ports:
- "10000:10000"
- "4040:4040"
environment:
SPARK_MODE: "master"
command: >
bash -c "
pip install delta-spark &&
/opt/bitnami/spark/sbin/start-thriftserver.sh
--conf spark.sql.extensions=io.delta.sql.DeltaSparkSessionExtension
--conf spark.sql.catalog.spark_catalog=org.apache.spark.sql.delta.catalog.DeltaCatalog
--conf spark.hadoop.fs.s3a.endpoint=http://minio:9000
--conf spark.hadoop.fs.s3a.access.key=minioadmin
--conf spark.hadoop.fs.s3a.secret.key=minioadmin123
--conf spark.hadoop.fs.s3a.path.style.access=true
--conf spark.sql.warehouse.dir=s3a://datalake/warehouse/
--hiveconf hive.server2.thrift.port=10000
--hiveconf hive.server2.thrift.bind.host=0.0.0.0 &&
tail -f /dev/null
"
volumes:
minio_data:
docker compose up -d
docker compose logs -f spark-thrift | grep "ThriftBinaryCLIService"
Шаг 2: Создание MinIO bucket¶
pip install minio
python3 - <<'EOF'
from minio import Minio
client = Minio("localhost:9000",
access_key="minioadmin",
secret_key="minioadmin123",
secure=False)
for bucket in ["datalake", "bronze", "silver", "gold"]:
if not client.bucket_exists(bucket):
client.make_bucket(bucket)
print(f"Created bucket: {bucket}")
EOF
Шаг 3: Инициализация dbt-проекта¶
pip install "dbt-spark[PyHive]" pyspark==3.5.0 delta-spark==3.1.0
dbt init analytics_lakehouse
Which database would you like to use?
[1] spark
Enter a number: 1
Desired project name: analytics_lakehouse
profile: analytics_lakehouse
Шаг 4: Написание profiles.yml¶
# ~/.dbt/profiles.yml
analytics_lakehouse:
target: dev
outputs:
dev:
type: spark
method: session
schema: "dbt_{{ env_var('USER', 'dev') }}"
threads: 2
server_side_parameters:
spark.master: "local[4]"
spark.driver.memory: "4g"
spark.sql.shuffle.partitions: "4"
spark.sql.adaptive.enabled: "true"
spark.jars.packages: >
io.delta:delta-core_2.12:2.4.0,
org.apache.hadoop:hadoop-aws:3.3.4
spark.sql.extensions: >
io.delta.sql.DeltaSparkSessionExtension
spark.sql.catalog.spark_catalog: >
org.apache.spark.sql.delta.catalog.DeltaCatalog
spark.hadoop.fs.s3a.endpoint: "http://localhost:9000"
spark.hadoop.fs.s3a.access.key: "{{ env_var('MINIO_ACCESS_KEY', 'minioadmin') }}"
spark.hadoop.fs.s3a.secret.key: "{{ env_var('MINIO_SECRET_KEY', 'minioadmin123') }}"
spark.hadoop.fs.s3a.path.style.access: "true"
spark.hadoop.fs.s3a.impl: >
org.apache.hadoop.fs.s3a.S3AFileSystem
thrift_local:
type: spark
method: thrift
host: localhost
port: 10000
user: "{{ env_var('USER', 'spark') }}"
schema: dbt_thrift_dev
auth: NONE
connect_timeout: 30
connect_retries: 3
threads: 2
server_side_parameters:
spark.sql.shuffle.partitions: "4"
Шаг 5: Проверка подключения¶
dbt debug --target dev
15:30:11 Connection:
15:30:11 host: None (session mode)
15:30:11 schema: dbt_yourusername
15:30:11 method: session
15:30:14 Connection test: [OK connection ok]
15:30:14 All checks passed!
dbt debug --target thrift_local
15:31:05 Connection:
15:31:05 host: localhost
15:31:05 port: 10000
15:31:05 schema: dbt_thrift_dev
15:31:05 auth: NONE
15:31:07 Connection test: [OK connection ok]
15:31:07 All checks passed!
Шаг 6: Тестовая модель¶
-- models/test_connection.sql
{{ config(materialized='view') }}
SELECT
1 AS test_value,
'connection ok' AS status,
CURRENT_TIMESTAMP() AS check_ts,
spark_version() AS spark_version
dbt run --select test_connection --target dev
dbt run --select test_connection --target thrift_local
Шаг 7: Сравнение результатов¶
dbt run --target dev 2>&1 | grep "Completed\|Error\|Running"
dbt run --target thrift_local 2>&1 | grep "Completed\|Error\|Running"
В session-режиме вы не увидите Spark UI (нет веб-сервера по умолчанию). В thrift-режиме откройте http://localhost:4040 - там будут видны все Job'ы, выполненные dbt.
Антипаттерны¶
Антипаттерн 1: Пароли в Git¶
prod:
password: "my_real_password"
token: "dapi0abc123"
Как только это попало в историю Git - учётные данные скомпрометированы, даже после удаления в следующем коммите. Всегда env_var().
Антипаттерн 2: Один shared thrift-кластер для всех разработчиков без изоляции¶
# Все разработчики используют один schema
dev_alice:
schema: analytics
dev_bob:
schema: analytics
Модели Alice перезаписывают модели Bob. Используйте динамический schema:
dev:
schema: "dbt_{{ env_var('USER') }}"
Теперь у каждого разработчика свой namespace: dbt_alice, dbt_bob.
Антипаттерн 3: Одни и те же Spark-конфигурации для dev и prod¶
dev:
server_side_parameters:
spark.sql.shuffle.partitions: "2000"
На dev с тестовыми датасетами в 1 МБ это создаёт 2000 пустых задач. На dev используйте минимальные конфигурации (shuffle.partitions: "4"), на prod - production-grade.
Антипаттерн 4: threads выше числа независимых веток в DAG¶
prod:
threads: 16
Если ваш DAG линейный (A → B → C → D), threads=16 ничего не ускорит: в каждый момент выполняется только одна модель. Проверьте DAG через dbt ls --output tree - если мало параллельных веток, увеличение threads бесполезно.
Антипаттерн 5: Отсутствие connect_retries для production¶
prod:
connect_timeout: 10
connect_retries: 0
Если Thrift Server кратковременно недоступен (рестарт, перегрузка), dbt немедленно упадёт. Добавьте connect_retries: 3 и connect_timeout: 60 - это обеспечивает устойчивость к кратковременным сбоям.
Антипаттерн 6: Использование session для production данных¶
prod:
method: session
server_side_parameters:
spark.master: "local[*]"
Session-режим работает на одной машине. Для терабайтных production-данных это неприемлемо. Session - только для dev и CI.
Чеклист: настройка dbt-spark¶
profiles.yml:
- Нет паролей и токенов в plain text - только
env_var() - Разные
schemaдля dev и prod (изоляция) schemaвключает имя пользователя для dev:dbt_{{ env_var('USER') }}- Разные
server_side_parameters.spark.sql.shuffle.partitionsдля dev (4–8) и prod (800–2000) connect_retries: 3иconnect_timeout: 60для prod thrift-подключенийthreadsсоответствует числу параллельных веток в DAG
Для session-метода:
spark.master: "local[*]"илиlocal[N]- Минимальные shuffle.partitions (4–8)
- Delta/Iceberg расширения подключены через
spark.sql.extensions - S3A/MinIO конфигурации для работы с объектным хранилищем
Для thrift-метода:
- Правильный
auth(NONE для dev без аутентификации, CUSTOM с паролем) use_ssl: trueесли трафик идёт через публичную сеть- production
shuffle.partitionsи AQE настройки
Для http-метода:
tokenтолько черезenv_var()use_ssl: trueвсегда- Правильный
endpoint(путь к SQL warehouse)
Домашнее задание¶
Дан файл profiles.yml с ошибками безопасности и конфигурации:
# profiles.yml - найдите и исправьте все ошибки
my_analytics:
target: production
outputs:
production:
type: spark
method: thrift
host: spark-cluster-prod.company.internal
port: 10000
user: analytics_user
password: "Str0ngPassw0rd!2024"
schema: PRODUCTION_SCHEMA
auth: NONE
connect_timeout: 5
threads: 100
server_side_parameters:
spark.sql.shuffle.partitions: "200"
spark.hadoop.fs.s3a.secret.key: "wJalrXUtnFEMI/K7MDENG/bPxRfiCYEXAMPLEKEY"
development:
type: spark
method: session
schema: PRODUCTION_SCHEMA
threads: 16
server_side_parameters:
spark.master: "local[*]"
spark.sql.shuffle.partitions: "2000"
spark.hadoop.fs.s3a.secret.key: "wJalrXUtnFEMI/K7MDENG/bPxRfiCYEXAMPLEKEY"
Задание 1. Найдите все проблемы в profiles.yml выше. Перечислите их с кратким объяснением, почему каждая является проблемой.
Задание 2. Перепишите profiles.yml корректно:
- Перенесите все секреты в переменные окружения
- Исправьте все конфигурационные ошибки
- Добавьте изоляцию схем для dev
- Настройте оптимальные Spark-параметры для каждого окружения
Задание 3. Запустите dbt debug с исправленным профилем против Docker Compose из лабораторной практики. Приложите вывод команды (с замаскированными чувствительными данными).
Задание 4. Создайте тестовую модель, которая выводит:
- Текущий таргет (
{{ target.name }}) - Текущую схему (
{{ target.schema }}) - Версию Spark (
spark_version()) - Количество shuffle-партиций (через
spark.conf.get(...)илиSET spark.sql.shuffle.partitions)
Запустите модель в dev и thrift_local таргетах. Что изменилось в выводе?