dbt vs PySpark: матрица выбора по сложности логики и объёму данных
dbt vs PySpark: матрица выбора по сложности логики и объёму данных
Концептуальный баттл: декларативность против императивности¶
Почему выбор инструмента - это архитектурное решение¶
В мире 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 требует:
- Сериализации данных из JVM в Python (через PyArrow или Pickle)
- Выполнения Python-кода в отдельном процессе
- Десериализации результата обратно в 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 - общие слова без обоснования.
Критерии оценки¶
- Точность выбора: верно ли определён квадрант матрицы?
- Обоснование: есть ли конкретные аргументы (объём, сложность, команда)?
- Риски: указаны ли потенциальные проблемы выбранного решения?
- Гибридность: правильно ли разделены ответственности в гибридных случаях?
- Практичность: реалистична ли предложенная архитектура для данной команды?