Dynamic Partition Pruning: механизм star-schema и условия срабатывания
Как Spark динамически отсекает ненужные партиции fact-таблицы на основе runtime-фильтров из dimension-таблиц, почему это работает только при BroadcastHashJoin, и как читать следы DPP в explain-плане.
В предыдущих уроках мы разобрали, как AQE адаптирует число shuffle-партиций после того, как данные уже записаны. Dynamic Partition Pruning работает иначе - он предотвращает само чтение ненужных данных ещё до того, как начнётся обработка. В этом его принципиальная ценность: вместо того чтобы оптимизировать уже считанные данные, он сокращает объём чтения на этапе scan.
Этот урок разбирает механизм DPP сверху вниз: от проблемы, которую он решает, через архитектурный контекст star-schema, до физических операторов в execution-плане и условий, при которых оптимизация срабатывает или отключается.
1. Жизнь до DPP: статическое vs динамическое отсечение¶
Статическое отсечение партиций (Static Partition Pruning)¶
Самый простой вид оптимизации, который Spark умеет делать с момента появления партиционирования. Если в SQL-запросе есть явный литеральный фильтр по колонке партиционирования - Spark читает только нужные папки в хранилище и полностью игнорирует остальные.
# Таблица sales партиционирована по колонке event_date
# Например: s3a://data/sales/event_date=2026-01-01/...
# s3a://data/sales/event_date=2026-01-02/...
# ... и так далее
# Запрос с литеральным фильтром
result = spark.table("sales").filter("event_date = '2026-01-01'")
result.explain(mode="simple")
# == Physical Plan ==
# FileScan parquet [amount, region_id, event_date]
# PartitionFilters: [isnotnull(event_date#12), (event_date#12 = 2026-01-01)]
# PushedFilters: []
# ReadSchema: struct<amount:double,region_id:int>
Spark видит event_date = '2026-01-01' на стадии построения логического плана - ещё до начала выполнения - и знает ровно одну папку, которую нужно открыть. Остальные 364 папки не прочитаются вообще. Это называется Predicate Pushdown применительно к партиционированию: условие «проталкивается» прямо в scan-оператор.
Ключевое свойство Static Partition Pruning: значение '2026-01-01' - константа, известная во время планирования. Catalyst видит литерал и сразу же выстраивает план с pruning.
Проблема: когда значение неизвестно при планировании¶
Типичный сценарий в реальных пайплайнах - фильтр приходит не в виде константы, а вычисляется из другой таблицы. Например, нужно найти продажи в те дни, когда проходили промоакции:
SELECT s.region_id, sum(s.amount) AS total
FROM sales s
JOIN promotions p ON s.event_date = p.promo_date
WHERE p.campaign = 'black_friday_2026'
Или то же самое на PySpark:
promos = spark.table("promotions").filter("campaign = 'black_friday_2026'")
result = spark.table("sales").join(promos, "event_date").groupBy("region_id").sum("amount")
Здесь Catalyst не может заранее знать, какие значения event_date попадут из promotions после фильтрации. Набор дат акции - это результат выполнения запроса, который станет известен только в runtime. Если партиционирование не учитывается - Spark вынужден читать все партиции sales, а потом уже отфильтровывать строки, которые не попали в join.
При таблице sales размером 10 ТБ и 3 годах истории это означает полный scan 1095 партиций, из которых реально нужны 5. Разница - в сотни раз по объёму I/O.
Dynamic Partition Pruning решает именно эту задачу: он позволяет Spark вычислить нужные значения ключа партиционирования в runtime, перед тем как начать scan fact-таблицы, и передать их в виде динамического фильтра прямо в файловую систему.
2. Архитектура Star-Schema: почему DPP создавался именно для неё¶
Модель «Звезда»¶
Star-schema - классическая модель данных в Data Warehouse и Lakehouse, где выделяют два типа таблиц:
Fact-таблица - большая, быстро растущая, содержит измеримые события: продажи, клики, транзакции, события логов. Миллиарды строк, десятки или сотни гигабайт, партиционирована по часто используемым измерениям (дата, регион, категория).
Dimension-таблицы - маленькие справочники, описывающие сущности: магазины, продукты, кампании, регионы. Тысячи или десятки тысяч строк, умещаются целиком в memory одного Executor.
Типичный аналитический запрос в star-schema выглядит так: применить фильтр к dimension, затем join с fact. Именно в этой точке рождается DPP: фильтр на dimension известен сразу, а DPP позволяет его использовать для отсечения партиций fact ещё до scan.
Почему партиционирование fact-таблицы критично¶
В star-schema fact-таблицу всегда партиционируют по тем измерениям, по которым чаще всего идёт join с dimensions. Это не случайность - это архитектурное решение, которое делает DPP возможным. Если fact-таблица не партиционирована по join-ключу, DPP физически не может ничего отсечь: нет папок, которые можно было бы пропустить.
# Сохранение fact-таблицы с правильным партиционированием
fact_df.write \
.mode("overwrite") \
.partitionBy("event_date", "region_id") \ # по тем же ключам, по которым будут join
.saveAsTable("sales")
# В результате получаем структуру:
# s3a://warehouse/sales/
# event_date=2026-01-01/
# region_id=1/ part-00000.parquet ...
# region_id=2/ part-00001.parquet ...
# event_date=2026-01-02/
# region_id=1/ ...
# ... (1095 папок для 3 лет данных)
3. Анатомия Dynamic Partition Pruning¶
Что происходит под капотом¶
DPP - это оптимизация на уровне Catalyst: трансформация логического плана, которая вставляет специальный оператор DynamicPruningExpression в scan-ноду fact-таблицы. Этот оператор в runtime вычисляет набор допустимых значений ключа партиционирования и передаёт их в файловую систему до открытия первого файла.
Рассмотрим пошагово на примере:
# Запрос: суммарные продажи в регионах, где прошла кампания "black_friday"
promos = spark.table("promotions").filter(col("campaign") == "black_friday_2026")
result = (
spark.table("sales") # fact: 10 ТБ, partition by event_date
.join(promos, "event_date") # join по partition-колонке
.groupBy("region_id")
.agg(F.sum("amount").alias("total_amount"))
)
Catalyst строит оптимальный физический план в несколько шагов:
Двойная роль dimension-таблицы¶
Здесь есть нюанс, который не сразу очевиден. Dimension-таблица (promotions) используется дважды в одном плане:
- Как сторона BroadcastHashJoin - данные dimension-таблицы рассылаются на все Executor-ы для выполнения join.
- Как источник значений для динамического фильтра - те же данные использует
DynamicPruningExpressionчтобы сформировать список допустимых значений partition-ключа.
Это не двойное чтение: Spark строит SubqueryBroadcast - специальный оператор, который читает dimension-таблицу один раз, строит broadcast-хэш-таблицу и одновременно передаёт набор ключей в scan-оператор fact-таблицы. Dimension читается ровно один раз, результат используется в двух местах плана.
Subquery Broadcast Injection: детали алгоритма¶
Когда Catalyst принимает решение применить DPP, он выполняет следующую последовательность трансформаций в логическом плане:
Шаг 1 - Фильтрация dimension-стороны. Spark выполняет scan dimension-таблицы с применением всех предикатов (в нашем случае campaign = 'black_friday_2026'). Результат - небольшое множество строк.
Шаг 2 - Извлечение join-ключей. Из отфильтрованных строк dimension-таблицы извлекается колонка, по которой идёт join с fact-таблицей. Это и есть будущий набор допустимых partition-значений: {2026-11-27, 2026-11-28, 2026-11-29}.
Шаг 3 - Broadcast-рассылка на все Executor-ы. Полученный набор ключей (обычно это несколько десятков или сотен значений) рассылается через механизм broadcast на каждый Executor. Это происходит до начала scan fact-таблицы.
Шаг 4 - Инжекция в FileScan. Оператор FileScan для fact-таблицы получает список допустимых partition-значений и передаёт его в Directory Listing. Файловая система возвращает только пути к папкам, которые присутствуют в этом списке. Все остальные папки не открываются и не читаются вообще.
Шаг 5 - BroadcastHashJoin на отфильтрованных данных. Теперь join происходит только на том подмножестве fact-данных, которое действительно может совпасть с dimension-стороной. Процент совпадений близок к 100% - нет «мусорных» строк, которые не найдут пары.
4. Жёсткие условия срабатывания¶
DPP - не универсальная оптимизация. Она применяется только при выполнении всех условий одновременно. Нарушение любого из них полностью отключает DPP.
Условие 1: Fact-таблица партиционирована по join-ключу¶
Это фундаментальное требование. Если fact-таблица не партиционирована или партиционирована по другой колонке, Spark физически не может отсечь файлы: нет структуры папок, которую можно использовать для pruning.
# ✅ DPP сработает: sales партиционирован по той же колонке, что используется в join
fact_df.write.partitionBy("event_date").saveAsTable("sales")
result = spark.table("sales").join(dim_df.filter(...), "event_date")
# ❌ DPP не сработает: join по product_id, а партиционирование по event_date
fact_df.write.partitionBy("event_date").saveAsTable("sales")
result = spark.table("sales").join(products.filter(...), "product_id")
# Spark прочтёт все партиции event_date, потому что product_id не является
# partition-колонкой - нет файловой структуры для отсечения
# ❌ DPP не сработает: таблица вообще не партиционирована
fact_df.write.saveAsTable("sales_no_partitioning")
result = spark.table("sales_no_partitioning").join(dim_df.filter(...), "event_date")
Частный случай: если fact-таблица партиционирована по нескольким колонкам (event_date, region_id), DPP может применяться независимо для каждой из них, если в запросе есть соответствующие join с dimension-таблицами.
Условие 2: Join идёт именно по partition-колонке¶
Мало того что таблица должна быть партиционирована - join в запросе обязан использовать ровно ту колонку, по которой выполнено партиционирование. Если join идёт по другому ключу или по вычисляемому выражению - DPP не применяется.
# ✅ DPP сработает: join по event_date = partition-колонка
spark.table("sales").join(dim_dates.filter(...), "event_date")
# ❌ DPP не сработает: join по выражению, а не по колонке напрямую
from pyspark.sql.functions import date_trunc
spark.table("sales").join(
dim_dates.filter(...),
date_trunc("month", col("event_date")) == date_trunc("month", col("date"))
# Вычисляемое выражение - Catalyst не может свести его к partition-колонке
)
# ❌ DPP не сработает: join по transformed-ключу
spark.table("sales").join(
dim_dates.filter(...),
col("event_date").cast("string") == col("promo_date")
# Cast создаёт новое выражение, не совпадающее с partition-колонкой
)
Условие 3: Тип Join должен быть совместим с DPP¶
DPP работает только с теми типами join, при которых Spark может использовать dimension-сторону для фильтрации fact-стороны без потери строк результата.
| Тип Join | DPP работает? | Почему |
|---|---|---|
INNER JOIN |
✅ Да | Только совпадающие строки - pruning безопасен |
LEFT SEMI JOIN |
✅ Да (fact слева) | По семантике аналогичен INNER - строки fact без пары отбрасываются |
RIGHT OUTER JOIN |
✅ Да (fact справа) | Fact-строки без пары отбрасываются |
LEFT OUTER JOIN |
❌ Нет (fact слева) | Строки fact без пары сохраняются → нельзя отсекать партиции |
FULL OUTER JOIN |
❌ Нет | Обе стороны сохраняют строки без пары |
CROSS JOIN |
❌ Нет | Нет ключа join - нечего использовать для pruning |
Причина ограничения для LEFT OUTER JOIN: если fact-таблица слева, то все её строки должны присутствовать в результате - даже те, которым нет пары в dimension. Следовательно, нельзя отсекать партиции fact по значениям из dimension: мы потеряли бы строки, которые должны присутствовать с null на dimension-стороне.
Условие 4: Join-стратегия - BroadcastHashJoin¶
Это самое неочевидное условие и самая частая причина «молчаливого» отказа DPP. По умолчанию DPP требует, чтобы dimension-сторона join использовала стратегию BroadcastHashJoin (BHJ).
Причина проста: механизм DPP физически построен на broadcast. SubqueryBroadcast - это broadacst dimension-данных на все Executor-ы. Именно через этот broadcast передаётся набор ключей для pruning. Если dimension-таблица слишком большая для broadcast и Spark выбирает SortMergeJoin, DPP не может сработать, потому что нет broadcast-артефакта, который можно использовать для фильтрации.
# ✅ DPP сработает: promotions < autoBroadcastJoinThreshold → BHJ автоматически
spark.conf.set("spark.sql.autoBroadcastJoinThreshold", "100m")
# promotions = 200 КБ < 100 МБ → Catalyst выберет BHJ → DPP активируется
# ✅ DPP сработает: явный broadcast-хинт
from pyspark.sql.functions import broadcast
spark.table("sales").join(broadcast(promotions.filter(...)), "event_date")
# ❌ DPP не сработает: BHJ отключён
spark.conf.set("spark.sql.autoBroadcastJoinThreshold", "-1")
# Spark выберет SortMergeJoin даже для маленьких таблиц → DPP отключится
# ❌ DPP не сработает: принудительный SMJ-хинт
spark.table("sales").hint("merge").join(promotions.filter(...), "event_date")
Настройка порога broadcast и DPP:
# Оба параметра важны вместе
spark.conf.set("spark.sql.autoBroadcastJoinThreshold", str(100 * 1024 * 1024)) # 100 МБ
spark.conf.set("spark.sql.optimizer.dynamicPartitionPruning.enabled", "true") # по умолчанию true с Spark 3.0
# Дополнительная настройка: использовать статистику для оценки выгодности DPP
spark.conf.set("spark.sql.optimizer.dynamicPartitionPruning.useStats", "true")
# Если статистики нет, Catalyst может ошибочно решить, что DPP невыгоден
# Установка в true заставляет использовать доступные table stats
Условие 5: Настройки Spark¶
Все три параметра должны быть включены для надёжной работы DPP:
spark = SparkSession.builder \
.appName("DPP-Production") \
.config("spark.sql.optimizer.dynamicPartitionPruning.enabled", "true") \
# ↑ главный включатель DPP (по умолчанию true начиная с Spark 3.0)
.config("spark.sql.adaptive.enabled", "true") \
# ↑ AQE улучшает DPP: в runtime может переключить SMJ → BHJ
# и тогда DPP тоже активируется там, где статически не мог
.config("spark.sql.autoBroadcastJoinThreshold", str(100 * 1024 * 1024)) \
# ↑ порог, до которого dimension автоматически broadcastится
# По умолчанию 10 МБ - часто слишком мало для реальных dimensions
.getOrCreate()
5. Практика: читаем планы выполнения¶
Понять, сработал ли DPP, можно двумя способами: через explain() в коде и через Spark UI. Разберём оба.
Подготовка синтетических данных¶
Создадим минимальную star-schema: таблицу фактов fact_sales и dimension dim_regions.
from pyspark.sql import SparkSession
from pyspark.sql import functions as F
from pyspark.sql.types import StructType, StructField, LongType, IntegerType, DoubleType, StringType, DateType
import datetime
spark = SparkSession.builder \
.appName("DPP-Demo") \
.config("spark.sql.optimizer.dynamicPartitionPruning.enabled", "true") \
.config("spark.sql.adaptive.enabled", "true") \
.config("spark.sql.autoBroadcastJoinThreshold", str(50 * 1024 * 1024)) \
.config("spark.sql.warehouse.dir", "/tmp/spark-warehouse") \
.enableHiveSupport() \
.getOrCreate()
# ─── 1. Dimension-таблица: регионы (маленькая, 4 строки) ────────────────────
dim_regions = spark.createDataFrame(
[
(1, "North", "Russia", "active"),
(2, "South", "Russia", "active"),
(3, "East", "China", "active"),
(4, "West", "Germany", "inactive"),
],
schema=StructType([
StructField("region_id", IntegerType(), False),
StructField("region_name", StringType(), True),
StructField("country", StringType(), True),
StructField("status", StringType(), True),
])
)
# ─── 2. Fact-таблица: продажи (большая, партиционирована по region_id) ───────
# Генерируем 10 млн строк с равномерным распределением по 4 регионам
fact_df = (
spark.range(1, 10_000_001)
.withColumn("region_id", ((F.col("id") % 4) + 1).cast("int"))
.withColumn("amount", (F.rand(seed=42) * 10000).cast("double"))
.withColumn("product_id", ((F.col("id") % 100) + 1).cast("int"))
.withColumn("event_date", (
F.date_add(F.lit(datetime.date(2025, 1, 1)), (F.col("id") % 365).cast("int"))
))
)
# Сохраняем с партиционированием по region_id - ключ для DPP
fact_df.write \
.mode("overwrite") \
.partitionBy("region_id") \
.saveAsTable("fact_sales")
# Результат: 4 папки region_id=1, region_id=2, region_id=3, region_id=4
spark.sql("ANALYZE TABLE fact_sales COMPUTE STATISTICS FOR ALL COLUMNS")
# Собираем статистику - Catalyst будет принимать лучшие решения
После сохранения у нас: 10 млн строк, разбитых на 4 партиции по region_id. Данные region_id=4 ("West", "Germany", "inactive") тоже записаны - это важно для демонстрации pruning.
Запрос с DPP¶
# Фильтруем dimension: оставляем только АКТИВНЫЕ регионы в России
# region_id=1 (North, Russia, active) и region_id=2 (South, Russia, active)
active_ru_regions = dim_regions.filter(
(F.col("country") == "Russia") & (F.col("status") == "active")
)
# active_ru_regions содержит region_id IN (1, 2)
# Join с fact по partition-колонке region_id
result = (
spark.table("fact_sales")
.join(active_ru_regions, "region_id") # join по partition-колонке
.groupBy("region_name", "country")
.agg(
F.sum("amount").alias("total_amount"),
F.count("*").alias("order_count"),
F.avg("amount").alias("avg_amount"),
)
)
Чтение explain()-плана¶
result.explain(mode="formatted")
Вывод (ключевые секции):
== Physical Plan ==
AdaptiveSparkPlan isFinalPlan=false
+- HashAggregate(keys=[region_name#7, country#8], ...)
+- Exchange hashpartitioning(region_name#7, country#8, 200), ...)
+- HashAggregate(keys=[region_name#7, country#8], ...)
+- Project [region_name#7, country#8, amount#3]
+- BroadcastHashJoin [region_id#2], [region_id#0], Inner, ...
:- FileScan parquet default.fact_sales [region_id#2, amount#3, ...]
: PartitionFilters: [isnotnull(region_id#2),
: dynamicpruning#42 [region_id#2]]
: ^^^^^^^^^^^^^^^^^^^^^^^^^^
: ВОТ ОН - МАРКЕР УСПЕШНОГО DPP
: ReadSchema: struct<amount:double,product_id:int,event_date:date>
+- BroadcastExchange HashedRelationBroadcastMode(...), ...
+- Filter ((country#8 = Russia) AND (status#9 = active))
+- Scan OneRowRelation [region_id#0, region_name#7, ...]
(dimension data - сканируется один раз)
== Subqueries ==
Subquery:0 Hosting operator id = 1 for dynamicpruning#42
+- BroadcastExchange HashedRelationBroadcastMode(List(cast(input[0,int,false]...))
+- Filter ((country#8 = Russia) AND (status#9 = active))
+- Scan OneRowRelation [region_id#0, ...]
Разберём ключевые маркеры:
dynamicpruning#42 [region_id#2] в PartitionFilters - это и есть главный маркер успешного DPP. Число 42 - это внутренний ID subquery. [region_id#2] - колонка, по которой применяется pruning. Если эта строка есть в PartitionFilters - DPP сработал.
BroadcastHashJoin в физическом плане - DPP использует broadcast. Если бы здесь был SortMergeJoin, dynamicpruning# не появился бы.
Секция Subqueries - показывает, что именно вычисляется для получения значений pruning. Здесь видно, что Spark читает dimension с фильтром country=Russia AND status=active и результат используется как набор partition-ключей.
Запрос без DPP: намеренно ломаем условия¶
Теперь посмотрим, как план выглядит без DPP. Отключаем broadcast:
# Явно отключаем broadcast для этого запроса
spark.conf.set("spark.sql.autoBroadcastJoinThreshold", "-1")
result_no_dpp = (
spark.table("fact_sales")
.join(active_ru_regions, "region_id")
.groupBy("region_name", "country")
.agg(F.sum("amount").alias("total_amount"))
)
result_no_dpp.explain(mode="formatted")
Фрагмент плана:
+- SortMergeJoin [region_id#2], [region_id#0], Inner
:- Sort [region_id#2 ASC NULLS FIRST], ...
: +- Exchange hashpartitioning(region_id#2, 200), ...
: +- FileScan parquet default.fact_sales [region_id#2, amount#3, ...]
: PartitionFilters: [isnotnull(region_id#2)]
: ^^^^^^^^^^^^^^^^^^^^^^^^^^
: НЕТ dynamicpruning# - DPP не сработал!
: ReadSchema: struct<amount:double,...>
+- Sort [region_id#0 ASC NULLS FIRST], ...
+- Exchange hashpartitioning(region_id#0, 200), ...
+- Filter ((country#8 = Russia) AND (status#9 = active))
+- Scan OneRowRelation [region_id#0, ...]
Ключевые различия:
SortMergeJoinвместоBroadcastHashJoin- это причина отсутствия DPP- В
PartitionFiltersтолькоisnotnull(region_id#2)- нетdynamicpruning# - Секции
Subqueriesв плане нет вообще - никаких подзапросов для pruning - Оба датафрейма проходят через
Exchange hashpartitioning- shuffle обеих сторон
Spark прочитает все 4 партиции fact-таблицы, а потом отбросит 50% данных в join (регионы 3 и 4). С DPP - читаются только 2 партиции из 4, сразу.
# Возвращаем broadcast обратно
spark.conf.set("spark.sql.autoBroadcastJoinThreshold", str(50 * 1024 * 1024))
Второй anti-pattern: join не по partition-колонке¶
# join по product_id - НЕ partition-колонка
dim_products = spark.createDataFrame(
[(1, "electronics"), (2, "clothing"), (3, "food")],
["product_id", "category"]
)
# fact_sales партиционирован по region_id, но join идёт по product_id
result_wrong_join = (
spark.table("fact_sales")
.join(broadcast(dim_products.filter(F.col("category") == "electronics")), "product_id")
.groupBy("category")
.agg(F.sum("amount").alias("total"))
)
result_wrong_join.explain(mode="formatted")
# В FileScan будет: PartitionFilters: [isnotnull(region_id#2)]
# НЕТ dynamicpruning# - BHJ есть, но join не по partition-колонке → DPP не сработал
Здесь BroadcastHashJoin присутствует (потому что мы используем broadcast()), но DPP всё равно не работает: partition-колонка region_id и join-колонка product_id - разные вещи. Pruning по region_id на основе фильтра product_id невозможен без дополнительной статистической корреляции между ними.
Spark UI: визуальная проверка¶
В Spark UI после выполнения запроса откройте вкладку SQL / Details. Там видна DAG-визуализация физического плана. При работающем DPP:
- Узел
FileScanсодержитPartitionFilters: dynamicpruning#N - Метрика number of files read показывает меньшее число файлов по сравнению с полным scan
- Метрика size of files read показывает суммарный объём прочитанных данных
Для сравнения «до и после»:
# Без DPP: читаем все 4 партиции
spark.conf.set("spark.sql.optimizer.dynamicPartitionPruning.enabled", "false")
result.write.mode("overwrite").format("noop").save() # запустить без реальной записи
# Смотрим в Spark UI: size of files read ≈ весь объём fact_sales
# С DPP: читаем только 2 партиции из 4
spark.conf.set("spark.sql.optimizer.dynamicPartitionPruning.enabled", "true")
result.write.mode("overwrite").format("noop").save()
# Смотрим в Spark UI: size of files read ≈ 50% от объёма (только region_id=1,2)
6. DPP с форматами Modern Data Stack¶
Parquet: базовый уровень¶
DPP работает с любым файловым форматом, который поддерживает partition-based directory listing. Parquet - стандартный выбор: Spark читает только нужные папки и внутри них применяет ещё один уровень оптимизации - column statistics из Parquet footer (min/max значения) позволяют пропускать row groups.
Delta Lake: дополнительный уровень через Transaction Log¶
Delta Lake добавляет поверх DPP свой механизм - Data Skipping через Transaction Log. Вместо того чтобы делать directory listing для определения нужных файлов, Spark читает _delta_log/*.json, где хранится статистика каждого файла (min/max по всем колонкам). DPP определяет нужные partition-значения, а Data Skipping дополнительно отсекает файлы внутри партиции по статистике колонок.
# DPP + Delta Lake
fact_df.write \
.format("delta") \
.mode("overwrite") \
.partitionBy("region_id") \
.saveAsTable("fact_sales_delta")
# Дополнительно: OPTIMIZE + ZORDER улучшает Data Skipping внутри файлов
spark.sql("OPTIMIZE fact_sales_delta ZORDER BY (event_date, product_id)")
# После OPTIMIZE: крупные Parquet файлы с хорошей статистикой → ещё меньше чтения
Apache Iceberg: partition evolution и hidden partitioning¶
Iceberg поддерживает hidden partitioning - данные партиционированы по трансформированному значению колонки (например, по месяцу из event_date), а пользователь пишет фильтры по исходной колонке. DPP в Iceberg интегрируется через Iceberg Catalog: plan-файлы позволяют отсекать не только partition-папки, но и конкретные data-файлы.
# Iceberg с hidden partitioning
# Физически партиционируем по месяцам, но join по event_date работает для DPP
spark.sql("""
CREATE TABLE iceberg.sales (
id BIGINT, region_id INT, amount DOUBLE, event_date DATE
) USING iceberg
PARTITIONED BY (months(event_date), region_id)
""")
# DPP будет работать при join по event_date или region_id
7. Взаимодействие DPP с AQE¶
AQE (Adaptive Query Execution) и DPP - взаимодополняющие оптимизации, которые вместе дают эффект больший, чем каждая по отдельности.
AQE как «спасательная сетка» для DPP¶
Рассмотрим сценарий: статически Spark считает dimension-таблицу слишком большой для broadcast (её размер превышает autoBroadcastJoinThreshold). Поэтому при планировании выбирается SortMergeJoin и DPP не активируется.
С включённым AQE ситуация меняется: Spark начинает выполнение с SortMergeJoin, но после shuffle-write на dimension-стороне собирает реальную статистику. Если выясняется, что после фильтрации dimension-данных стало мало - AQE переключает план на BroadcastHashJoin. И в этот момент DPP тоже активируется - уже в runtime.
Именно поэтому включение AQE является «условием 5» для DPP: с AQE значительно больше запросов получают DPP, потому что runtime-переключение на BHJ открывает возможность для pruning там, где статическое планирование её не видело.
AdaptiveSparkPlan в explain()¶
При включённом AQE explain() показывает план в состоянии isFinalPlan=false до выполнения и isFinalPlan=true после. DPP в explain() виден только в финальном плане (после выполнения), если AQE переключил стратегию:
# Перед выполнением
result.explain(mode="formatted")
# AdaptiveSparkPlan isFinalPlan=false ← статический план
# Закешировать и посмотреть финальный план
result.cache()
result.count() # trigger execution
result.explain(mode="formatted")
# AdaptiveSparkPlan isFinalPlan=true ← финальный план после AQE-переключений
# Здесь dynamicpruning# уже видно если DPP сработал
8. Anti-patterns и диагностика¶
Anti-pattern 1: партиционирование по high-cardinality колонке¶
# ❌ Партиционирование по user_id с кардинальностью 100 миллионов
fact_df.write.partitionBy("user_id").saveAsTable("fact_events")
# Результат: 100 млн папок в хранилище → directory listing занимает минуты
# Object Storage (S3, MinIO) имеет лимиты на list-операции
# ✅ Партиционирование по колонке с разумной кардинальностью
fact_df.write.partitionBy("event_date", "region_id").saveAsTable("fact_events")
# event_date: ~1000 уникальных дат за 3 года (разумно)
# region_id: ~100 регионов (разумно)
# Вместе: ~100 000 партиций - это уже на грани, но управляемо
Практическое правило: кардинальность partition-колонки должна быть в диапазоне от 10 до 10 000 уникальных значений. При больших значениях directory listing становится дороже самого pruning.
Anti-pattern 2: слепое доверие к DPP без проверки explain¶
Включённый конфиг dynamicPartitionPruning.enabled=true не гарантирует применение оптимизации. DPP молча не срабатывает при нарушении условий. Всегда проверяйте:
def check_dpp(df):
"""Проверить, сработал ли DPP в физическом плане."""
plan_str = df._sc._jvm.PythonSQLUtils.explainString(
df._jdf.queryExecution(), "formatted"
)
if "dynamicpruning#" in plan_str:
print("✅ DPP активирован")
else:
print("❌ DPP НЕ сработал - проверьте условия:")
print(" 1. Fact-таблица партиционирована по join-ключу?")
print(" 2. Join по partition-колонке (не по выражению)?")
print(" 3. Тип join совместим (INNER, LEFT SEMI, RIGHT OUTER)?")
print(" 4. Dimension достаточно мала для BroadcastHashJoin?")
print(" 5. spark.sql.optimizer.dynamicPartitionPruning.enabled = true?")
check_dpp(result)
Anti-pattern 3: UDF в join-предикате¶
from pyspark.sql.functions import udf
@udf("string")
def normalize_region(r):
return r.strip().lower()
# ❌ UDF в join-условии: Catalyst не может выразить это через partition-фильтр
result = spark.table("fact_sales").join(
dim_regions,
normalize_region(col("region_name")) == normalize_region(col("dim_region_name"))
)
# DPP не сработает: UDF - чёрный ящик для Catalyst, нет связи с partition-колонкой
# ✅ Нормализация отдельно, до join
dim_clean = dim_regions.withColumn("region_id_clean", col("region_id"))
spark.table("fact_sales").join(dim_clean.filter(...), col("region_id") == col("region_id_clean"))
Anti-pattern 4: join через subquery без partition-ключа¶
# ❌ Subquery возвращает не partition-ключ, а вычисляемое значение
subquery = spark.sql("""
SELECT DISTINCT month(promo_date) AS promo_month FROM promotions WHERE campaign='BF'
""")
result = spark.table("fact_sales").join(
subquery,
F.month("event_date") == F.col("promo_month")
)
# DPP не сработает: `month(event_date)` - это выражение, не partition-колонка event_date
# ✅ Оставаться на уровне partition-ключа
promo_dates = spark.sql("SELECT DISTINCT promo_date FROM promotions WHERE campaign='BF'")
result = spark.table("fact_sales").join(promo_dates, "event_date")
# join по event_date = partition-колонка → DPP сработает
9. Best Practices и чек-лист¶
Чек-лист дата-инженера¶
Перед написанием запроса с join на fact-таблицу проверьте:
Проектирование схемы:
- Fact-таблица партиционирована по колонкам, по которым чаще всего идут join с dimensions
- Кардинальность partition-колонки от 10 до ~10 000 уникальных значений
- Dimension-таблицы намеренно держатся маленькими (< 100–500 МБ)
- Собрана статистика (
ANALYZE TABLE ... COMPUTE STATISTICS)
Написание запроса:
- Join идёт по partition-колонке без трансформаций (
event_date, неmonth(event_date)) - Тип join совместим с DPP: INNER, LEFT SEMI, RIGHT OUTER
- Нет UDF в join-условиях и partition-фильтрах
- Для очень маленьких dimensions - явный
broadcast()хинт для надёжности
Настройки Spark:
spark.conf.set("spark.sql.optimizer.dynamicPartitionPruning.enabled", "true") # уже true по умолчанию
spark.conf.set("spark.sql.adaptive.enabled", "true") # уже true в Spark 3.2+
spark.conf.set("spark.sql.autoBroadcastJoinThreshold", str(100 * 1024 * 1024)) # увеличить с 10МБ до 100МБ
spark.conf.set("spark.sql.optimizer.dynamicPartitionPruning.useStats", "true") # использовать статистику
Верификация:
- Запустить
result.explain(mode="formatted")и найтиdynamicpruning#в PartitionFilters - В Spark UI сравнить
size of files readс ожидаемым (должно быть пропорционально доле нужных партиций) - Убедиться что в плане
BroadcastHashJoin, а неSortMergeJoin
Итоговая сводка: когда DPP работает¶
Домашнее задание¶
Дан следующий код:
# Таблицы
orders = spark.table("orders") # 500 ГБ, партиционирована по order_date
stores = spark.table("stores") # 2 МБ, 300 магазинов
products = spark.table("products") # 150 МБ, 200 000 товаров
# Запрос
result = (
orders
.join(stores.filter(col("country") == "RU"), col("store_id") == col("id"))
.join(products.filter(col("category") == "electronics"), "product_id")
.filter(col("order_date") >= "2026-01-01")
.groupBy("country", "category")
.agg(F.sum("amount").alias("total"))
)
Задание 1. Определите, какие join в этом запросе потенциально могут использовать DPP. Для каждого join объясните: по какой колонке идёт join, является ли эта колонка partition-колонкой orders, и подходит ли dimension-таблица для broadcast.
Задание 2. Перепишите запрос так, чтобы максимизировать шанс срабатывания DPP. Добавьте явные broadcast() хинты где нужно, вынесите filter(order_date >= ...) на правильное место, переупорядочьте join-ы если необходимо.
Задание 3. Запустите оба варианта (исходный и оптимизированный) и сравните explain(mode="formatted"). Найдите dynamicpruning# в оптимизированном варианте. Запишите, сколько Exchange узлов (shuffle) присутствует в каждом варианте.
Задание 4. Предположим, что products (150 МБ) не помещается в broadcast при дефолтном пороге 10 МБ. Как можно решить проблему тремя разными способами - через настройку, через hint и через переработку схемы данных?
В следующем уроке разберём Bucketing - как предварительно партиционировать данные по хэшу ключа при записи, чтобы полностью устранить shuffle при повторяющихся join по тому же ключу.