dbt vs PySpark: матрица выбора по сложности логики и объёму данных

dbt vs PySpark: матрица выбора по сложности логики и объёму данных

platform

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

Почему выбор инструмента - это архитектурное решение

В мире Data Engineering существует соблазн найти «лучший инструмент» и использовать его везде. PySpark-инженеры пишут на PySpark всё - от ingestion до финальных витрин. Команды, увлечённые dbt, пытаются вместить в SQL-макросы логику, для которой SQL не предназначен.

Оба подхода ведут к техническому долгу: в первом случае - тысячи строк нечитаемого PySpark-кода там, где достаточно десяти строк SQL; во втором - чудовищные Jinja-макросы, имитирующие рекурсию, которой нет в реляционной алгебре.

Выбор между dbt и PySpark - это не вопрос вкуса. Это архитектурное решение, которое влияет на:

  • Maintainability: сможет ли новый аналитик понять пайплайн через год?
  • Performance: какой инструмент лучше оптимизирует конкретный тип операций?
  • Governance: как обеспечить тестирование, документацию, lineage?
  • Hiring: какие навыки нужны команде?
  • Cost: время разработки × зарплата + время выполнения × стоимость кластера

Опытный дата-инженер должен владеть обоими инструментами и знать, когда использовать каждый. Это и есть тема урока.

Парадигма dbt: «Что мы хотим получить»

dbt реализует декларативный подход: вы описываете желаемый результат в виде SELECT-запроса, а dbt и Spark самостоятельно определяют, как его получить. Вы не указываете, как перемешивать данные, как их кэшировать, как разбивать на партиции - это решает Catalyst Optimizer.

-- dbt: декларативно описываем желаемую таблицу
{{ config(materialized='incremental', file_format='delta') }}

SELECT
    user_id,
    event_date,
    COUNT(*)                AS session_count,
    SUM(revenue)            AS total_revenue,
    AVG(session_duration)   AS avg_duration
FROM {{ ref('stg_events') }}
GROUP BY user_id, event_date

Catalyst получает этот SQL и строит физический план: выбирает HashAggregate или SortAggregate, решает, применять ли BroadcastHashJoin, оптимизирует порядок предикатов. Инженер не управляет этим напрямую.

Сила декларативности: читаемость, лаконичность, возможность для SQL-аналитиков писать трансформации без знания Spark internals. Витрина из 50 JOIN-ов записывается в 100 строк SQL - понятных любому аналитику.

Ограничение декларативности: SQL - это реляционная алгебра. Есть класс задач, которые принципиально нельзя эффективно выразить в реляционной алгебре: рекурсия, итеративные алгоритмы, обработка бинарных данных, интеграция с Python ML-библиотеками.

Парадигма PySpark: «Как мы хотим это сделать»

PySpark реализует императивный подход: вы описываете последовательность шагов - трансформаций - которые Spark должен выполнить. Полный контроль над графом вычислений.

# PySpark: императивно управляем каждым шагом
from pyspark.sql import functions as F
from pyspark.sql.window import Window

# Явно указываем порядок операций
df = spark.read.parquet("s3a://bronze/events/")

# Явно управляем кэшированием промежуточных результатов
df_cleaned = (
    df
    .filter(F.col("event_type").isNotNull())
    .withColumn("event_date", F.to_date("event_ts"))
    .persist(StorageLevel.MEMORY_AND_DISK_SER)  # Явное кэширование
)

# Явно управляем партиционированием перед тяжёлой операцией
df_repartitioned = df_cleaned.repartition(200, "user_id")

window_spec = Window.partitionBy("user_id").orderBy("event_ts")
df_sessions = df_repartitioned.withColumn(
    "session_id",
    F.sum(
        F.when(
            F.col("event_ts").cast("long") -
            F.lag("event_ts").over(window_spec).cast("long") > 1800,
            1
        ).otherwise(0)
    ).over(window_spec)
)

df_sessions.write.partitionBy("event_date").parquet("s3a://silver/sessions/")

Сила императивности: полный контроль. Можно написать Python-функцию любой сложности, вызвать любую библиотеку, управлять памятью, реализовать рекурсию через .foreach(), интегрировать ML-модели.

Ограничение императивности: многословность и сложность поддержки. Трансформация, которая в SQL занимает 10 строк, в PySpark - 40-50. При сотнях таблиц это накапливается в тысячи строк кода, где трудно разобраться без автора.

Точка пересечения: оба используют Spark SQL Engine

Важное понимание, которое часто упускают: и dbt, и PySpark в конечном счёте выполняются на одном движке - Spark SQL Engine с Catalyst Optimizer. dbt генерирует SQL и отправляет его через Thrift Server. PySpark DataFrame API компилируется в логический план, который затем оптимизируется Catalyst.

Это означает:

  • Один и тот же SELECT в dbt и spark.sql("SELECT ...") в PySpark даст одинаковый физический план
  • DataFrame API (.groupBy().agg()) и эквивалентный SQL часто дают одинаковый план
  • Python UDF (udf()) - принципиальное исключение: они выполняются в Python-процессе и не оптимизируются Catalyst

Миф «PySpark всегда быстрее dbt» - неверен. Правильнее: «PySpark UDF медленнее SQL-функций; PySpark DataFrame API часто эквивалентен SQL по производительности; PySpark даёт больше контроля там, где Catalyst не может автоматически оптимизировать».


Координата №1: объём данных и физическая оптимизация

Диапазон, где dbt работает отлично

dbt со Spark SQL движком уверенно справляется с объёмами от гигабайтов до нескольких десятков терабайтов. В этом диапазоне:

  • Catalyst Optimizer эффективно оптимизирует планы: выбирает тип join, применяет predicate pushdown, оптимизирует порядок операций
  • AQE (Adaptive Query Execution) динамически подстраивает число shuffle-партиций, обнаруживает и исправляет skew
  • Delta Lake / Iceberg предоставляют Z-ordering, compaction и file skipping, которые dbt использует через tblproperties

Для типичной аналитической платформы с 10-100 GB ежедневного прироста данных dbt в связке со Spark - оптимальный выбор. Управление физическим хранением через partition_by, clustered_by, location_root (см. урок 5) покрывает большинство сценариев оптимизации.

Где начинается зона PySpark

Существует несколько сигналов, что объём или характер данных требует PySpark:

Петабайтный масштаб и ручной тюнинг памяти. При экстремальных объёмах дефолтных настроек AQE недостаточно. Нужно управлять:

# PySpark: тонкий тюнинг под конкретный джоб
spark = (
    SparkSession.builder
    .config("spark.executor.memory", "32g")
    .config("spark.executor.memoryOverhead", "8g")   # Off-heap для shuffle
    .config("spark.memory.fraction", "0.8")           # 80% heap под Spark
    .config("spark.memory.storageFraction", "0.3")    # 30% из этих 80% под кэш
    .config("spark.sql.shuffle.partitions", "4000")   # Вручную под объём
    .config("spark.reducer.maxSizeInFlight", "96m")   # Буфер reduce-фазы
    .getOrCreate()
)

В dbt эти параметры можно передать через server_side_parameters в profiles.yml, но они применяются глобально для всей сессии, а не для конкретной тяжёлой операции.

Экстремальный Data Skew, который не исправляет AQE. AQE справляется с умеренным skew (в несколько раз). Когда один ключ содержит 80% данных - нужен ручной salting:

# PySpark: ручной salting для перекошенного join
import random

# Добавляем случайный суффикс к горячему ключу
df_skewed = df_orders.withColumn(
    "user_id_salted",
    F.when(
        F.col("user_id") == "power_user_12345",  # Горячий ключ
        F.concat(F.col("user_id"), F.lit("_"), (F.rand() * 10).cast("int").cast("string"))
    ).otherwise(F.col("user_id"))
)

# Дублируем справочник под каждый salt-суффикс
df_users_exploded = df_users.withColumn(
    "user_id_salted",
    F.explode(F.array([
        F.concat(F.col("user_id"), F.lit(f"_{i}")) for i in range(10)
    ]))
)

# Теперь join без skew
result = df_skewed.join(df_users_exploded, "user_id_salted")

В dbt это невозможно выразить в SQL - salting требует процедурного кода. Можно написать макрос, но читаемость будет ужасной.

Итеративные алгоритмы с кэшированием промежуточных состояний. Некоторые алгоритмы требуют многократного чтения одного датафрейма:

# PySpark: итеративный PageRank - не выразимо в SQL
def pagerank(edges, num_iterations=10, damping=0.85):
    vertices = edges.select("src").union(edges.select("dst")).distinct()
    ranks = vertices.withColumn("rank", F.lit(1.0))

    for i in range(num_iterations):
        # Кэшируем перед итерацией - иначе план растёт линейно
        ranks = ranks.persist(StorageLevel.MEMORY_AND_DISK)

        contributions = edges.join(ranks, edges.src == ranks.id) \
            .groupBy("dst") \
            .agg(F.sum(F.col("rank") / F.col("out_degree")).alias("contribution"))

        ranks = contributions.withColumn(
            "rank",
            F.lit(1 - damping) + F.lit(damping) * F.col("contribution")
        )
        ranks.unpersist()

    return ranks

Без .persist() Spark пересчитывает ranks заново на каждой итерации. В dbt нет механизма кэширования промежуточных результатов между операциями.

Работа с нереляционными данными на Bronze-уровне. dbt ожидает таблицы - строки и столбцы. Бинарные файлы (Avro, Protobuf, ORC с кастомными SerDe), изображения, аудио, видео - всё это требует предварительной обработки на PySpark.


Координата №2: сложность бизнес-логики и типы данных

Слепая зона dbt: где SQL превращается в ад

Рекурсивные и итеративные алгоритмы. SQL поддерживает рекурсивные CTE (WITH RECURSIVE), но Spark SQL - нет (на момент написания). Даже если бы поддерживал, рекурсивный CTE в SQL неэффективен для глубоких деревьев - он плохо параллелизуется.

Пример: построение сессий из потока событий. В SQL - либо оконные функции (работают только для простых случаев), либо сложные self-join (экспоненциально дорогие). В PySpark - элегантная итерация:

# PySpark: построение сессий через Gap Detection
from pyspark.sql.window import Window
from pyspark.sql import functions as F

SESSION_GAP = 30 * 60  # 30 минут

window = Window.partitionBy("user_id").orderBy("event_ts")

df_sessions = (
    df_events
    .withColumn(
        "prev_event_ts",
        F.lag("event_ts").over(window)
    )
    .withColumn(
        "is_new_session",
        F.when(
            F.col("event_ts").cast("long") -
            F.col("prev_event_ts").cast("long") > SESSION_GAP,
            1
        ).when(F.col("prev_event_ts").isNull(), 1)
        .otherwise(0)
    )
    .withColumn(
        "session_id",
        F.sum("is_new_session").over(
            window.rowsBetween(Window.unboundedPreceding, Window.currentRow)
        )
    )
    .withColumn("global_session_id", F.concat_ws("_", "user_id", "session_id"))
)

В dbt это можно сделать через оконные функции SQL, но только для простых случаев. Как только сессии требуют динамического порога, нескольких типов событий-прерывателей или вложенных состояний - SQL становится нечитаемым.

Обработка глубоко вложенных JSON. dbt имеет ограниченную поддержку get_json_object, from_json и explode. Но динамические JSON-структуры неизвестной вложенности требуют рекурсивных Python-функций:

# PySpark: рекурсивный парсинг произвольно вложенного JSON
def flatten_schema(schema, prefix=""):
    """Рекурсивно разворачивает любой StructType в flat список колонок."""
    fields = []
    for field in schema.fields:
        field_name = f"{prefix}.{field.name}" if prefix else field.name
        if isinstance(field.dataType, StructType):
            fields.extend(flatten_schema(field.dataType, field_name))
        elif isinstance(field.dataType, ArrayType):
            fields.append(field_name)
        else:
            fields.append(field_name)
    return fields

# Применяем к датафрейму с произвольной вложенностью
raw = spark.read.json("s3a://bronze/api_responses/")
flat_columns = flatten_schema(raw.schema)
df_flat = raw.select([F.col(c).alias(c.replace(".", "_")) for c in flat_columns])

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

Интеграция с ML-экосистемой. Feature Engineering для ML часто требует:

  • Векторизации (VectorAssembler)
  • Нормализации (StandardScaler, MinMaxScaler)
  • Категориального кодирования (StringIndexer, OneHotEncoder)
  • Применения сохранённых ML-моделей (PipelineModel.load())

Всё это недоступно в dbt SQL. PySpark MLlib предоставляет нативный API:

from pyspark.ml.feature import VectorAssembler, StandardScaler
from pyspark.ml import Pipeline
from pyspark.ml.regression import GBTRegressor

# Feature pipeline - невозможно в dbt
assembler = VectorAssembler(
    inputCols=["age", "ltv", "days_since_last_order", "total_orders"],
    outputCol="features"
)
scaler = StandardScaler(inputCol="features", outputCol="scaled_features")
model = GBTRegressor(featuresCol="scaled_features", labelCol="churn_label")

pipeline = Pipeline(stages=[assembler, scaler, model])
fitted = pipeline.fit(df_train)
predictions = fitted.transform(df_test)

Pandas UDF для сложных бизнес-функций. Когда алгоритм нельзя векторизовать SQL-функциями, но нужна высокая производительность - Pandas UDF (Vectorized UDF):

import pandas as pd
from pyspark.sql.functions import pandas_udf

@pandas_udf("double")
def calculate_geohash_distance(lat1: pd.Series, lon1: pd.Series,
                               lat2: pd.Series, lon2: pd.Series) -> pd.Series:
    """Haversine distance через NumPy - быстрее Python UDF в 10-100 раз."""
    import numpy as np
    R = 6371.0  # Радиус Земли в км
    dlat = np.radians(lat2 - lat1)
    dlon = np.radians(lon2 - lon1)
    a = np.sin(dlat/2)**2 + np.cos(np.radians(lat1)) * np.cos(np.radians(lat2)) * np.sin(dlon/2)**2
    return R * 2 * np.arcsin(np.sqrt(a))

# Использование в pipeline
df_delivery = df_orders.withColumn(
    "distance_km",
    calculate_geohash_distance("pickup_lat", "pickup_lon", "delivery_lat", "delivery_lon")
)

Pandas UDF выполняется через Apache Arrow - намного быстрее обычных Python UDF, но всё равно медленнее JVM-функций. В dbt доступны только Spark SQL built-in функции.

Слепая зона PySpark: где Python плодит технический долг

Классический реляционный ETL. Когда нужно построить 200 аналитических витрин - star schema из fact/dimension таблиц - PySpark превращается в кошмар поддержки:

# PySpark: 50 строк кода для того, что в dbt занимает 10
def build_fct_orders(spark, env):
    df_orders = spark.table(f"{env}.stg_orders")
    df_order_items = spark.table(f"{env}.stg_order_items")
    df_users = spark.table(f"{env}.dim_users")
    df_products = spark.table(f"{env}.dim_products")

    df_items_agg = (
        df_order_items
        .groupBy("order_id")
        .agg(
            F.sum("quantity").alias("total_items"),
            F.sum("line_total").alias("subtotal"),
            F.sum(F.col("quantity") * F.col("discount_per_unit")).alias("total_discount")
        )
    )

    df_result = (
        df_orders
        .join(df_items_agg, "order_id", "left")
        .join(df_users.select("user_id", "country", "segment"), "user_id", "left")
        .withColumn("net_amount", F.col("subtotal") - F.col("total_discount"))
        .withColumn("event_date", F.to_date("created_at"))
        .select(
            "order_id", "user_id", "country", "segment",
            "total_items", "subtotal", "total_discount", "net_amount",
            "status", "created_at", "event_date"
        )
    )

    (
        df_result
        .write
        .mode("overwrite")
        .partitionBy("event_date")
        .format("delta")
        .saveAsTable(f"{env}.fct_orders")
    )

Тот же результат в dbt:

-- dbt: 10 строк SQL, читаемых любым аналитиком
SELECT
    o.order_id,
    o.user_id,
    u.country,
    u.segment,
    SUM(i.quantity)                                AS total_items,
    SUM(i.line_total)                              AS subtotal,
    SUM(i.quantity * i.discount_per_unit)          AS total_discount,
    SUM(i.line_total) - SUM(i.quantity * i.discount_per_unit) AS net_amount,
    o.status,
    o.created_at,
    date(o.created_at)                             AS event_date
FROM {{ ref('stg_orders') }} o
LEFT JOIN {{ ref('stg_order_items') }} i USING (order_id)
LEFT JOIN {{ ref('dim_users') }} u USING (user_id)
GROUP BY o.order_id, o.user_id, u.country, u.segment, o.status, o.created_at

Data Lineage и документация. В PySpark lineage нужно строить вручную или через сторонние инструменты (OpenLineage, Spline, Marquez). Это дополнительная работа, которую часто не делают - и через год никто не знает, откуда берётся та или иная колонка в витрине.

В dbt lineage строится автоматически через ref(). Функция dbt docs generate создаёт полноценный портал документации с графом зависимостей. Тест relationships проверяет ссылочную целостность. Всё это «из коробки», без дополнительного кода.

Тестирование. PySpark-пайплайны тестируют через pytest + chispa (или pytest-spark). Это требует написания тестовых данных, мок-SparkSession, ручных assertions. Для 200 витрин - несколько тысяч строк тестового кода.

В dbt generic тесты (not_null, unique, accepted_values) подключаются одной строкой в YAML. Сотни проверок - десятки строк конфигурации.

# dbt: 5 строк YAML = 5 тестов
columns:
  - name: order_id
    tests: [not_null, unique]
  - name: status
    tests: [{accepted_values: {values: ['pending', 'completed', 'cancelled']}}]
# PySpark: 30+ строк Python = 2 теста
def test_fct_orders_order_id_not_null(spark, df_fct_orders):
    null_count = df_fct_orders.filter(F.col("order_id").isNull()).count()
    assert null_count == 0, f"Found {null_count} null order_ids"

def test_fct_orders_order_id_unique(spark, df_fct_orders):
    total = df_fct_orders.count()
    distinct = df_fct_orders.select("order_id").distinct().count()
    assert total == distinct, f"Found {total - distinct} duplicate order_ids"

Catalyst Optimizer: почему декларативный SQL часто быстрее ручного кода

Как Catalyst оптимизирует SQL

Когда dbt отправляет SQL на Spark, он проходит через несколько фаз оптимизации:

SQL Text
  → Parsing (AST)
    → Analysis (resolve refs, types)
      → Logical Optimization
        → Physical Planning
          → Code Generation (Tungsten)

На этапе Logical Optimization Catalyst применяет десятки правил:

  • Predicate Pushdown: WHERE event_date = '2024-01-01' перемещается максимально близко к источнику данных, ограничивая объём читаемых данных
  • Column Pruning: читаются только колонки, которые реально нужны в запросе
  • Join Reordering: Catalyst может переставлять join-ы чтобы уменьшить промежуточные результаты
  • Constant Folding: 100 * 3 вычисляется на этапе оптимизации, не в рантайме
  • Subquery Elimination: дублирующиеся подзапросы объединяются

На этапе Physical Planning Catalyst выбирает конкретные алгоритмы:

  • BroadcastHashJoin для маленьких таблиц (< spark.sql.autoBroadcastJoinThreshold)
  • SortMergeJoin для больших таблиц
  • HashAggregate или SortAggregate в зависимости от кардинальности

Почему ручной PySpark DataFrame API не всегда умнее

Парадоксально, но ручной DataFrame API иногда мешает Catalyst оптимизировать:

# ПЛОХО: явный repartition ломает оптимизацию join
df_orders = spark.table("orders").repartition(400)  # Принудительный shuffle
df_users = spark.table("users")

result = df_orders.join(df_users, "user_id")
# Catalyst хотел BroadcastHashJoin (users маленькая),
# но принудительный repartition до join вставил Exchange →
# теперь оба shuffle, BroadcastHashJoin невозможен
# ЛУЧШЕ: позволяем Catalyst выбрать
result = spark.table("orders").join(spark.table("users"), "user_id")
# Catalyst видит, что users < 10 MB → BroadcastHashJoin без shuffle

В dbt, поскольку вы пишете SQL, у Catalyst полная свобода оптимизации. Явные вмешательства отсутствуют.

Реальное узкое место: Python UDF

Единственное место, где PySpark принципиально медленнее SQL - это Python UDF. Каждый вызов Python UDF требует:

  1. Сериализации данных из JVM в Python (через PyArrow или Pickle)
  2. Выполнения Python-кода в отдельном процессе
  3. Десериализации результата обратно в JVM

Это в 10-100 раз медленнее JVM-функций. В dbt вы не можете вызвать Python UDF - только встроенные SQL-функции, которые выполняются нативно в JVM.

# МЕДЛЕННО: Python UDF
@F.udf("string")
def normalize_phone(phone):                         # Python execution
    import re
    return re.sub(r'[^0-9+]', '', str(phone) or '')

df = df.withColumn("phone_clean", normalize_phone("phone"))

# БЫСТРО: SQL-функции (JVM native)
df = df.withColumn(
    "phone_clean",
    F.regexp_replace(F.col("phone").cast("string"), r'[^0-9+]', '')
)

Вывод: если задачу можно решить встроенными SQL-функциями - dbt и PySpark-SQL одинаково быстры. Если нужен Python - используйте Pandas UDF (векторизованный), а не обычный Python UDF.


Матрица выбора: квадрант решений

Четыре квадранта

Квадрант 1: малый/средний объём + простая логика → чистый dbt

Характеристики:

  • Данные: до 100 GB на таблицу, реляционные источники (PostgreSQL, MySQL, Kafka через Landing)
  • Логика: JOIN, GROUP BY, оконные функции, CASE WHEN
  • Команда: аналитики + analytics engineers, знающие SQL
  • Примеры: CRM-отчётность, финансовые витрины, маркетинговые дашборды, продуктовые метрики

Почему dbt: скорость разработки максимальная. Новая витрина - несколько часов вместо дней. Аналитики могут писать модели сами без участия инженеров. Документация, тесты, lineage - всё из коробки.

Риски: через год проект вырастет, и некоторые модели попадут в квадранты 2 или 4. Это нормально - рефакторинг конкретных моделей, не всего проекта.

Квадрант 2: большой объём + простая логика → dbt + Spark-конфигурации

Характеристики:

  • Данные: 100 GB - несколько TB на таблицу, партиционированные источники
  • Логика: та же реляционная, но на больших объёмах
  • Примеры: E-commerce с 10M+ заказов, телеком с миллиардами событий, банковская транзакционность

Почему dbt с конфигурациями: SQL-логика не изменилась - только объём. Решение: правильные partition_by, clustered_by, server_side_parameters, строгие contract. AQE + правильная физическая организация данных покрывает большинство сценариев.

{{ config(
    file_format='delta',
    partition_by={'field': 'event_date', 'data_type': 'date'},
    clustered_by=['user_id'],
    buckets=512,
    pre_hook="SET spark.sql.shuffle.partitions = 1200"
) }}

Когда переходить в квадрант 4: когда появляется data skew, который AQE не исправляет, или когда shuffle занимает доминирующее время в джобе.

Квадрант 3: малый/средний объём + сложная логика → PySpark / Pandas

Характеристики:

  • Данные: до 100 GB, но сложный формат или алгоритм
  • Логика: ML feature engineering, нестандартный парсинг, рекурсия, iterative algorithms
  • Примеры: R&D, прототипы ML, обработка бинарных форматов, графовые алгоритмы

Почему PySpark: гибкость Python-экосистемы важнее декларативности. На малом объёме производительность не критична - важна скорость прототипирования.

Типичный путь: прототип на PySpark → проверка идеи → оптимизация горячих путей → при необходимости перевод части логики в dbt

Квадрант 4: большой объём + сложная логика → чистый PySpark

Характеристики:

  • Данные: TB - PB, часто нереляционные форматы
  • Логика: custom алгоритмы, ML, сложный CDC, streaming
  • Примеры: рекомендательные системы на миллиардах событий, real-time антифрод, IoT телеметрия

Почему PySpark: единственный вариант. Нужен ручной тюнинг под специфику джоба, Python-библиотеки, низкоуровневый контроль.

Критично: такие пайплайны сложно поддерживать. Инвестируйте в документацию (OpenLineage), тестирование (chispa, pytest-spark) и code review.

Детализированная матрица сигналов

Сигнал dbt PySpark
Трансформация - чистый SQL (JOIN/AGG/WINDOW)
Нужен Python-код внутри трансформации
Команда - SQL-аналитики
Команда - Senior Data Engineers
Важна документация «из коробки»
Нужен кастомный lineage через OpenLineage
Data Skew, который AQE не исправляет
Iterative / recursive algorithms
Работа с MLlib / sklearn / NumPy
Бинарные форматы (Avro/Protobuf)
Объём до 100 GB ✓ (предпочтительно)
Объём > 1 TB, простая логика ✓ с конфигурациями
Объём > 1 TB, сложная логика
Нужны тесты «из коробки»
Частые изменения бизнес-логики
Стабильная алгоритмическая логика

Гибридная архитектура: лучшее из двух миров

Рекомендуемый паттерн для modern data platform

Большинство production data platform в итоге приходят к гибридной архитектуре: PySpark для Bronze (сырые данные), dbt для Silver/Gold (аналитика). Это не компромисс - это правильное разделение инструментов по их сильным сторонам.

[Источники]
    │
    ▼
[Bronze Layer]  ← PySpark
    • Сырая загрузка из Kafka, S3, APIs
    • Парсинг нестандартных форматов (Avro, Protobuf)
    • Deduplication CDC-событий
    • Первичная нормализация схем
    │
    ▼
[Silver Layer]  ← dbt
    • Очистка и стандартизация
    • Реляционная нормализация
    • Generic тесты на все ключевые поля
    │
    ▼
[Gold Layer]  ← dbt
    • Аналитические витрины
    • Dimensional modeling (fact/dim)
    • Data Contracts, строгие гарантии
    • Документация и lineage
    │
    ▼
[ML/Streaming]  ← PySpark (отдельная ветка)
    • Feature engineering
    • Model inference
    • Structured Streaming

Пример: Bronze-загрузчик на PySpark

# Bronze: PySpark для сложного ingestion
# Задача: загрузить Avro-события из Kafka в Delta Lake

from pyspark.sql import SparkSession
from pyspark.sql import functions as F
from pyspark.sql.avro.functions import from_avro

spark = SparkSession.builder \
    .config("spark.sql.extensions", "io.delta.sql.DeltaSparkSessionExtension") \
    .getOrCreate()

# Читаем схему из Schema Registry
with open("schemas/user_event.avsc") as f:
    avro_schema = f.read()

df_kafka = (
    spark.readStream
    .format("kafka")
    .option("kafka.bootstrap.servers", "kafka:9092")
    .option("subscribe", "user_events")
    .load()
)

# Десериализация Avro - невозможно в dbt SQL
df_decoded = (
    df_kafka
    .withColumn("event", from_avro(F.col("value"), avro_schema))
    .select("event.*", F.col("timestamp").alias("_kafka_timestamp"))
)

# Добавляем ingestion metadata
df_bronze = (
    df_decoded
    .withColumn("_ingested_at", F.current_timestamp())
    .withColumn("_source_offset", F.col("offset"))
    .withColumn("_source_partition", F.col("partition"))
)

# Пишем в Delta с merge для идемпотентности
(
    df_bronze.writeStream
    .format("delta")
    .outputMode("append")
    .option("checkpointLocation", "s3a://checkpoints/bronze_events/")
    .table("bronze.raw_user_events")
)

После того как Bronze-таблица наполнена - dbt берёт управление:

-- Silver: dbt модель поверх Bronze
-- stg_user_events.sql
{{ config(materialized='incremental', file_format='delta') }}

SELECT
    event_id,
    user_id,
    event_type,
    CAST(properties:session_id   AS STRING) AS session_id,
    CAST(properties:page_url     AS STRING) AS page_url,
    CAST(properties:revenue      AS DECIMAL(18,2)) AS revenue,
    _ingested_at,
    date(_ingested_at)                       AS event_date
FROM {{ source('bronze', 'raw_user_events') }}
{% if is_incremental() %}
WHERE _ingested_at > (SELECT MAX(_ingested_at) FROM {{ this }})
{% endif %}

PySpark сделал то, что не умеет dbt (Kafka + Avro). dbt делает то, что не умеет PySpark без боли (тесты, документация, lineage, понятный SQL).

Точки интеграции между PySpark и dbt

Интеграция происходит через общий Data Lakehouse (Delta/Iceberg). PySpark пишет в Delta-таблицу → dbt читает её через source(). Никакой специальной связи не нужно - только договорённость о схеме (data contract для Bronze).

# sources.yml: dbt "видит" таблицу, написанную PySpark
sources:
  - name: bronze
    tables:
      - name: raw_user_events
        description: Kafka → Avro → Delta через PySpark Streaming
        loaded_at_field: _ingested_at
        freshness:
          warn_after: {count: 5, period: minutes}
          error_after: {count: 30, period: minutes}
        columns:
          - name: event_id
            tests: [not_null, unique]
          - name: _ingested_at
            tests: [not_null]

Кейс-батл: два реальных пайплайна

Кейс A: Маркетинговая воронка атрибуции - выбор dbt

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

Вводные данные:

  • 50 таблиц-источников: Google Ads, Facebook, email, SEO, push-уведомления
  • Объём: 10-50 GB в сутки
  • Логика часто меняется: маркетологи регулярно хотят другую модель атрибуции (last-touch, linear, time-decay)
  • Команда: 3 analytics engineers, 5 аналитиков-маркетологов

Выбор: dbt

Обоснование:

  • 50 источников = 50 staging-моделей + несколько intermediate + несколько Gold-витрин. В dbt это 60-70 SQL-файлов с автоматическим lineage. В PySpark - 60-70 Python-файлов с ручным управлением зависимостями.
  • Логика часто меняется: SQL-модели маркетологи могут читать и частично изменять сами. PySpark-код для них непонятен.
  • Разные модели атрибуции легко реализуются как dbt-макросы или отдельные модели:
-- models/marts/marketing/fct_attribution_linear.sql
-- Linear attribution: равный вес всем touchpoint'ам
{{ config(materialized='table') }}

WITH touchpoints AS (
    SELECT
        user_id,
        conversion_id,
        channel,
        COUNT(*) OVER (PARTITION BY conversion_id) AS total_touches
    FROM {{ ref('int_marketing_touchpoints') }}
)

SELECT
    conversion_id,
    channel,
    SUM(1.0 / total_touches)  AS attributed_conversions,
    SUM(revenue / total_touches) AS attributed_revenue
FROM touchpoints
GROUP BY conversion_id, channel
-- models/marts/marketing/fct_attribution_last_touch.sql
-- Last-touch: весь credit последнему каналу
{{ config(materialized='table') }}

WITH ranked AS (
    SELECT *,
        ROW_NUMBER() OVER (
            PARTITION BY conversion_id ORDER BY touched_at DESC
        ) AS rn
    FROM {{ ref('int_marketing_touchpoints') }}
)

SELECT
    conversion_id,
    channel,
    1                AS attributed_conversions,
    revenue          AS attributed_revenue
FROM ranked
WHERE rn = 1

Маркетологи могут переключаться между моделями атрибуции, меняя одну строку в дашборде. Два SQL-файла против одного сложного PySpark-скрипта с условиями.

Кейс B: IoT телеметрия с аномалиями - выбор PySpark

Задача: промышленное предприятие собирает телеметрию с 10 000 датчиков. Данные приходят в бинарном формате (кастомный Protobuf). Нужно: декодировать бинарные данные, вычислить скользящие окна аномалий, дедуплицировать по сессиям оборудования.

Вводные данные:

  • Источник: 10 000 датчиков → 100 000 событий/сек → Kafka
  • Объём: 5 TB/день
  • Формат: бинарный Protobuf с кастомной схемой, версионируется
  • Логика: аномалия = значение датчика отклоняется от скользящего среднего за 5 минут более чем на 3σ
  • Команда: 2 data engineers с ML-бэкграундом

Выбор: PySpark

Обоснование:

  • Protobuf-декодирование требует Python-библиотеки (protobuf). В dbt нет механизма подключить Python-библиотеку для обработки данных.
  • Скользящее среднее за 5 минут × 10 000 датчиков на 5 TB - это оконная функция по физическому времени (RANGE BETWEEN), которая в Spark SQL работает, но требует тонкой настройки partitioning под 10 000 групп.
  • Дедупликация по сессиям - stateful operation, которую лучше делать через Spark Structured Streaming с watermark.
# PySpark: парсинг Protobuf
import sensor_pb2  # Скомпилированный Protobuf

def decode_sensor_event(binary_data):
    """Декодирование кастомного Protobuf - только в Python."""
    event = sensor_pb2.SensorEvent()
    event.ParseFromString(bytes(binary_data))
    return (
        event.sensor_id,
        event.timestamp_us,
        event.value,
        event.quality_flag,
        event.firmware_version
    )

decode_udf = F.udf(decode_sensor_event, StructType([...]))

# Streaming pipeline
df_raw = (
    spark.readStream.format("kafka")
    .option("subscribe", "sensor_telemetry")
    .load()
)

df_decoded = df_raw.withColumn(
    "parsed",
    decode_udf(F.col("value"))
).select("parsed.*", "_kafka_timestamp")

# Скользящее среднее и аномалии
window_spec = Window.partitionBy("sensor_id") \
    .orderBy(F.col("timestamp_us").cast("long")) \
    .rangeBetween(-5 * 60 * 1_000_000, 0)  # 5 минут в микросекундах

df_anomalies = (
    df_decoded
    .withColumn("rolling_mean", F.avg("value").over(window_spec))
    .withColumn("rolling_std", F.stddev("value").over(window_spec))
    .withColumn(
        "is_anomaly",
        F.abs(F.col("value") - F.col("rolling_mean")) > 3 * F.col("rolling_std")
    )
    .filter(F.col("is_anomaly"))
)

df_anomalies.writeStream.format("delta").table("bronze.sensor_anomalies")

dbt здесь неприменим принципиально: Protobuf-декодирование требует Python.


Сравнение пайплайнов: explain plan в двух вариантах

Посмотрим на физические планы одинаковой трансформации в dbt и PySpark:

-- dbt (компилируется в SQL):
SELECT user_id, SUM(amount) AS total_revenue
FROM stg_orders
WHERE status = 'completed'
GROUP BY user_id
# PySpark DataFrame API:
result = (
    spark.table("stg_orders")
    .filter(F.col("status") == "completed")
    .groupBy("user_id")
    .agg(F.sum("amount").alias("total_revenue"))
)

Физический план - идентичен:

== Physical Plan ==
*(2) HashAggregate(keys=[user_id#1], functions=[sum(amount#2)])
+- Exchange hashpartitioning(user_id#1, 200), ENSURE_REQUIREMENTS
   +- *(1) HashAggregate(keys=[user_id#1], functions=[partial_sum(amount#2)])
      +- *(1) Filter (isnotnull(status#3) AND (status#3 = completed))
         +- *(1) ColumnarToRow
            +- FileScan parquet [user_id#1,amount#2,status#3] ...
               PartitionFilters: [], DataFilters: [isnotnull(status#3), (status#3 = completed)]
               PushedFilters: [IsNotNull(status), EqualTo(status,completed)]

Catalyst применил PushedFilters одинаково для SQL и DataFrame API. Производительность идентична.


Anti-patterns: как не нужно делать

Антипаттерн 1: всё в PySpark (включая простые витрины)

# ПЛОХО: PySpark для простой витрины - 60 строк
def build_fct_daily_revenue(spark, date):
    df = spark.table("stg_orders") \
        .filter(F.col("order_date") == date) \
        .join(spark.table("dim_users").select("user_id", "country"), "user_id") \
        .join(spark.table("dim_products").select("product_id", "category"), "product_id") \
        .groupBy("country", "category", "order_date") \
        .agg(
            F.sum("revenue").alias("total_revenue"),
            F.count("order_id").alias("order_count"),
            F.countDistinct("user_id").alias("unique_buyers")
        ) \
        .withColumn("avg_order_value", F.col("total_revenue") / F.col("order_count"))
    # Где тесты? Где документация? Где lineage?
    df.write.mode("overwrite").saveAsTable("gold.fct_daily_revenue")

Результат через год: 200 таких функций, никто не знает откуда берутся цифры, новый аналитик тратит неделю чтобы разобраться в одной витрине.

Антипаттерн 2: сложная логика в SQL-макросах dbt

{# ПЛОХО: попытка сделать рекурсию через Jinja-macro #}
{% macro build_hierarchy(node_id, depth=0, max_depth=10) %}
    {% if depth < max_depth %}
        SELECT {{ node_id }} as id, {{ depth }} as level
        UNION ALL
        {{ build_hierarchy(
            "(SELECT parent_id FROM hierarchy WHERE id = " ~ node_id ~ ")",
            depth + 1,
            max_depth
        ) }}
    {% endif %}
{% endmacro %}

SQL без WITH RECURSIVE не умеет в рекурсию. Этот макрос развернётся в безумный UNION ALL глубиной 10 уровней - нечитаемый, медленный, ограниченный по глубине.

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

Антипаттерн 3: Python UDF вместо SQL-функций там, где они есть

# ПЛОХО: Python UDF для того, что есть в SQL
@F.udf("string")
def extract_year(date_str):
    return str(date_str)[:4] if date_str else None

# ХОРОШО: SQL-функция (100x быстрее)
F.year(F.col("created_at")).cast("string")

Всегда проверяйте, есть ли встроенная SQL-функция, прежде чем писать UDF.

Антипаттерн 4: использование dbt для streaming/real-time

dbt - это batch-инструмент. Инкрементальные модели - это micro-batch с задержкой в минуты. Если бизнес требует секундной latency - нужен PySpark Structured Streaming или Flink.

# ХОРОШО: Spark Structured Streaming для real-time
df_stream = (
    spark.readStream
    .format("kafka")
    .option("subscribe", "orders")
    .load()
    .writeStream
    .trigger(processingTime="10 seconds")  # Micro-batch каждые 10 секунд
    .format("delta")
    .table("gold.realtime_orders")
)

Миграция: от legacy PySpark ETL к dbt

Стратегия постепенной миграции

Команды с legacy PySpark-кодом часто хотят перейти на dbt, но не могут сделать это одним шагом. Рекомендуемая стратегия:

Фаза 1: Параллельная работа. Существующий PySpark-пайплайн продолжает работать. Начинаем писать dbt-модели поверх Bronze-таблиц, которые создаёт PySpark. Постепенно переносим Silver/Gold логику в dbt.

Фаза 2: Инкрементальный перенос. Переносим в dbt всё, что является чистым SQL (JOIN, AGG). Оставляем в PySpark то, что требует Python (parsin, ML, streaming).

Фаза 3: Четкое разделение. PySpark отвечает за Bronze (ingestion). dbt отвечает за Silver и Gold. Граница - source() в dbt.

# До миграции: весь пайплайн в одном PySpark-скрипте
def full_etl_pipeline():
    # Bronze
    df_raw = load_and_parse_avro("s3a://kafka-sink/orders/")
    df_bronze = deduplicate_cdc(df_raw)
    df_bronze.write.saveAsTable("bronze.raw_orders")

    # Silver - ПЕРЕНЕСЁМ В dbt
    df_silver = (
        df_bronze
        .filter(F.col("order_id").isNotNull())
        .withColumn("order_date", F.to_date("created_at"))
    )
    df_silver.write.saveAsTable("silver.stg_orders")

    # Gold - ПЕРЕНЕСЁМ В dbt
    df_gold = (
        df_silver
        .groupBy("order_date", "country")
        .agg(F.sum("revenue").alias("daily_revenue"))
    )
    df_gold.write.saveAsTable("gold.fct_daily_revenue")

После миграции:

# PySpark отвечает только за Bronze
def bronze_ingestion():
    df_raw = load_and_parse_avro("s3a://kafka-sink/orders/")
    df_bronze = deduplicate_cdc(df_raw)
    df_bronze.write.saveAsTable("bronze.raw_orders")
    # Дальше - dbt
-- dbt: Silver
-- stg_orders.sql
SELECT order_id, user_id, revenue, status,
       TO_DATE(created_at) AS order_date
FROM {{ source('bronze', 'raw_orders') }}
WHERE order_id IS NOT NULL
-- dbt: Gold
-- fct_daily_revenue.sql
SELECT order_date, country, SUM(revenue) AS daily_revenue
FROM {{ ref('stg_orders') }}
JOIN {{ ref('dim_users') }} USING (user_id)
GROUP BY order_date, country

Домашнее задание: Architecture Decision Record

Контекст

Вы - senior data engineer в крупном ритейлере. Бизнес предъявил три требования. Вам нужно написать Architecture Decision Record (ADR) - короткий технический документ с обоснованием архитектурного выбора.

Требование 1: Ежедневная маркетинговая отчётность

  • Источники: Google Analytics, Яндекс.Метрика, Facebook Ads API, внутренняя CRM - 8 таблиц
  • Объём: 500 MB - 2 GB в сутки
  • Команда: 2 analytics engineers, 4 бизнес-аналитика
  • SLA: витрины готовы к 9:00 утра
  • Метрики: ROAS, CAC, LTV по каналам, конверсионные воронки

Требование 2: Антифрод для платёжных транзакций

  • Источник: Kafka-топик платёжного шлюза, 2 000 транзакций/сек
  • Объём: 1 TB/день
  • Алгоритм: скользящее окно 30 секунд, если > 5 транзакций с одного IP - блокировка
  • Latency: решение < 100 мс, иначе транзакция пропускается
  • Команда: 3 data engineers с distributed systems бэкграундом

Требование 3: Персонализированные рекомендации

  • Источник: история просмотров (5 TB/день), покупки (200 GB/день), каталог товаров (10 GB)
  • Алгоритм: матричная факторизация (ALS) → топ-10 рекомендаций для каждого пользователя
  • Частота: пересчёт 1 раз в сутки ночью
  • Команда: 2 ML-инженера, 1 data engineer

Структура ADR для каждого требования

## ADR-{N}: {Название компонента}

### Статус
Предложено / Принято / Отклонено

### Контекст
[Описание задачи и ограничений]

### Решение
[Выбор: dbt / PySpark / Гибрид]

### Обоснование
[Конкретные аргументы: объём, сложность логики, команда, SLA]

### Альтернативы
[Что было рассмотрено и почему отклонено]

### Последствия
[Риски, требования к команде, инфраструктурные изменения]

### Метрики успеха
[Как поймём, что решение правильное]

Хороший ADR содержит конкретные цифры и отвечает на вопрос «почему не альтернатива». Плохой ADR - общие слова без обоснования.

Критерии оценки

  • Точность выбора: верно ли определён квадрант матрицы?
  • Обоснование: есть ли конкретные аргументы (объём, сложность, команда)?
  • Риски: указаны ли потенциальные проблемы выбранного решения?
  • Гибридность: правильно ли разделены ответственности в гибридных случаях?
  • Практичность: реалистична ли предложенная архитектура для данной команды?