dbt-spark адаптер: методы подключения thrift, http, session и profiles.yml

dbt-spark - адаптер, соединяющий dbt Core и Apache Spark. Разбираем методы thrift, http, session, анатомию profiles.yml, передачу Spark-конфигов и интеграцию с Delta Lake и MinIO.

platform

Архитектурный мост между 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 - логин/пароль через SASL
  • KERBEROS - 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.

Порядок поиска:

  1. Переменная окружения DBT_PROFILES_DIR - если задана, dbt ищет в этой директории
  2. Аргумент --profiles-dir при запуске команды
  3. Файл profiles.yml в текущей директории проекта
  4. ~/.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 ..."

Для диагностики медленных моделей:

  1. В dbt логах найдите, какая модель долго выполняется
  2. В Spark UI найдите соответствующий Job по времени или имени таблицы
  3. Проверьте метрики Stage'ей: Shuffle Read Size, Spill, Task Duration
  4. Скорректируйте 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 таргетах. Что изменилось в выводе?