Dynamic Partition Pruning: механизм star-schema и условия срабатывания

Как Spark динамически отсекает ненужные партиции fact-таблицы на основе runtime-фильтров из dimension-таблиц, почему это работает только при BroadcastHashJoin, и как читать следы DPP в explain-плане.

optimization

В предыдущих уроках мы разобрали, как 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) используется дважды в одном плане:

  1. Как сторона BroadcastHashJoin - данные dimension-таблицы рассылаются на все Executor-ы для выполнения join.
  2. Как источник значений для динамического фильтра - те же данные использует 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 по тому же ключу.