Bucketing: sort-merge join без shuffle - настройка и ограничения

Как бакетизация переносит стоимость shuffle с каждого запроса на момент записи данных, физический механизм Bucket Co-location в Sort-Merge Join, жёсткие условия срабатывания и архитектурные ограничения в современных Lakehouse.

optimization

В предыдущем уроке мы разобрали Dynamic Partition Pruning - механизм, который сокращает объём читаемых данных. Bucketing решает другую задачу: он устраняет самую дорогую операцию при join больших таблиц - Shuffle. Вместо того чтобы перемешивать данные по сети при каждом запросе, стоимость перемешивания переносится на момент записи таблицы. Написали один раз, join делаем тысячу раз - и каждый раз без Shuffle.


1. Анатомия проблемы: цена Shuffle в Sort-Merge Join

Три фазы классического Sort-Merge Join

Когда две большие таблицы не помещаются в память для Broadcast Hash Join, Spark использует Sort-Merge Join (SMJ). Этот алгоритм работает корректно и масштабируется на любые объёмы данных, но состоит из трёх последовательных фаз, каждая из которых стоит дорого.

Фаза 1 - Shuffle (самая дорогая). Каждая из двух таблиц полностью перераспределяется по сети так, чтобы все строки с одинаковым значением join-ключа оказались на одном Executor-е. Для этого Spark вычисляет hash(join_key) % numPartitions для каждой строки и отправляет её на соответствующий Executor. Все данные обоих датасетов проходят через сеть - даже те, которые в итоге не найдут пары.

Фаза 2 - Сортировка. После получения данных каждый Executor сортирует свою порцию обеих таблиц по join-ключу. Это необходимо для следующей фазы. Сортировка - O(N log N) и может сопровождаться Spill на диск при нехватке памяти.

Фаза 3 - Слияние (Merge). Два отсортированных итератора сливаются по алгоритму merge sort. Это линейная O(N+M) операция, самая дешёвая из трёх.

Фаза Shuffle - это:

  • Disk I/O на Map-стороне: каждый Map Task записывает shuffle-файлы на локальный диск Executor-а перед отправкой.
  • Сетевой трафик: все данные обеих таблиц передаются по TCP. При join 500 ГБ + 50 ГБ - 550 ГБ трафика каждый раз.
  • Disk I/O на Reduce-стороне: полученные блоки снова пишутся на диск перед сортировкой при нехватке RAM (Spill).
  • Latency: Spark не может начать Reduce-фазу, пока не завершатся все Map Tasks - это barrier synchronization.

И всё это происходит при каждом выполнении запроса. Если этот join запускается 50 раз в день в рамках BI-отчётов, то 550 ГБ × 50 = 27.5 ТБ сетевого трафика ежедневно только на одном join.

Идея Bucketing: перенести shuffle на момент записи

Bucketing - это предварительное распределение строк по N файлам (бакетам) на основе hash(bucket_key) % N. Выполняется один раз при записи таблицы. Если две таблицы забакетированы одинаково (по тому же ключу и с тем же числом бакетов), то строки с одинаковым customer_id уже находятся в файлах с одинаковым номером в обеих таблицах. Shuffle становится избыточным - данные уже распределены так, как нужно для join.

При чтении каждый Executor берёт бакет с одинаковым номером из обеих таблиц. Он точно знает, что все строки с одинаковым customer_id уже в его паре файлов - больше никуда ходить не нужно. Shuffle ликвидирован полностью.


2. Партиционирование vs Бакетизация

Оба механизма разбивают данные на части, но решают принципиально разные задачи. Путаница между ними - одна из самых распространённых ошибок при проектировании хранилищ.

Характеристика Партиционирование Бакетизация
Структура на диске Иерархия директорий (/dt=2026-01-01/) Файлы с номерами (bucket_00001.parquet)
Ключ разбиения Любая колонка с разумной кардинальностью Join-ключ с высокой кардинальностью
Цель Pruning: не читать лишние данные Устранение Shuffle при join
Кардинальность ключа Низкая (дни, регионы, категории) Высокая (user_id, order_id)
Число разделов = число уникальных значений ключа Фиксированное N (задаётся при записи)
Механизм оптимизации File path filtering при scan Hash co-location при join
Тип данных в Metastore Partition metadata (путь, значение) Bucket metadata (ключ, N, sort-ключ)
Требует Metastore? Нет (работает через read.parquet(path)) Да (только через saveAsTable + spark.table())

Партиционирование хорошо работает вместе с бакетизацией. Типичная архитектура: таблица партиционирована по event_date (для pruning) и забакетирована по user_id (для join). Spark сначала отсекает ненужные даты (Partition Pruning), потом читает только нужные бакеты без Shuffle (Bucket Join).

При запросе WHERE event_date = '2026-01-01' JOIN users ON user_id Spark:

  1. Открывает только директорию event_date=2026-01-01/ - Partition Pruning.
  2. Внутри неё читает только нужные бакеты - Bucket Co-location.
  3. Join выполняется без Shuffle - данные уже в нужном месте.

3. Физический механизм Bucket Co-location

Как Spark понимает, что Shuffle не нужен

Spark хранит метаданные о бакетизации в Hive Metastore (или совместимом каталоге). Когда Catalyst строит план выполнения join, он читает метаданные обеих таблиц и сравнивает:

  • Совпадают ли bucket-ключи с join-ключами запроса?
  • Одинаково ли число бакетов в обеих таблицах?
  • Одинакова ли функция хэширования?

Если все три условия выполнены - Catalyst помечает обе стороны join как BucketedScan и убирает узлы Exchange (hashpartitioning) из физического плана. Вместо Shuffle каждый Task напрямую открывает соответствующую пару файлов.

Что происходит с файлами на диске

# Представьте, что мы записали таблицу с bucketBy(4, "user_id")
# На диске появятся файлы:
# warehouse/orders/
#   part-00000-{uuid}_00000.c000.snappy.parquet   ← bucket 0
#   part-00000-{uuid}_00001.c000.snappy.parquet   ← bucket 1
#   part-00000-{uuid}_00002.c000.snappy.parquet   ← bucket 2
#   part-00000-{uuid}_00003.c000.snappy.parquet   ← bucket 3

# Суффикс _00000, _00001, ... - это номер бакета в имени файла.
# Spark использует его для сопоставления файлов при join.

Число файлов при записи = numBuckets × numPartitions (в других смыслах). Если у вас 200 Spark-задач записи и bucketBy(32, "user_id"), на диске появится 200 файлов, по 32 бакета, но несколько файлов с одинаковым номером бакета. Spark при чтении объединяет все файлы одного бакета в один логический «бакет».

Метаданные в Metastore

# Проверить метаданные бакетизации через SQL
spark.sql("DESCRIBE EXTENDED orders_bucketed").show(50, truncate=False)
# В выводе будет секция Table Information:
# Num Buckets:              32
# Bucket Columns:           [user_id]
# Sort Columns:             [user_id]
# ...

# Или через Python API
spark.sql("SHOW CREATE TABLE orders_bucketed").show(1, truncate=False)

4. Настройка и синтаксис

Запись бакетированной таблицы

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

spark = SparkSession.builder \
    .appName("Bucketing-Demo") \
    .config("spark.sql.warehouse.dir", "/tmp/spark-warehouse") \
    .enableHiveSupport() \
    .getOrCreate()

# ─── 1. Генерация тестовых данных ────────────────────────────────────────────
orders_df = (
    spark.range(1, 50_000_001)  # 50 млн заказов
    .withColumn("user_id",    (F.col("id") % 1_000_000).cast("long"))  # 1M уникальных пользователей
    .withColumn("amount",     (F.rand(seed=1) * 5000 + 10).cast("double"))
    .withColumn("status",     F.when(F.col("id") % 5 == 0, "cancelled").otherwise("completed"))
    .withColumn("order_date", F.date_add(F.lit("2025-01-01"), (F.col("id") % 365).cast("int")))
)

users_df = (
    spark.range(0, 1_000_000)  # 1M пользователей
    .withColumn("user_id",  F.col("id").cast("long"))
    .withColumn("country",  F.when(F.col("id") % 3 == 0, "RU")
                             .when(F.col("id") % 3 == 1, "US")
                             .otherwise("DE"))
    .withColumn("tier",     F.when(F.col("id") % 10 == 0, "premium").otherwise("standard"))
    .withColumn("email",    F.concat(F.lit("user"), F.col("id"), F.lit("@mail.com")))
    .drop("id")
)
# ─── 2. Запись orders с бакетизацией ─────────────────────────────────────────
(
    orders_df.write
    .mode("overwrite")
    .bucketBy(32, "user_id")    # 32 бакета по колонке user_id
    .sortBy("user_id")           # сортировка ВНУТРИ каждого бакета по тому же ключу
    .saveAsTable("orders_bucketed")
    # ВАЖНО: saveAsTable, а не write.parquet(path)!
    # Только saveAsTable регистрирует метаданные в Metastore.
)

# ─── 3. Запись users с теми же параметрами бакетизации ───────────────────────
(
    users_df.write
    .mode("overwrite")
    .bucketBy(32, "user_id")    # ТО ЖЕ число бакетов, ТОТ ЖЕ ключ
    .sortBy("user_id")           # ТА ЖЕ сортировка - убирает Sort-фазу из SMJ
    .saveAsTable("users_bucketed")
)

Разберём параметры детально:

bucketBy(numBuckets, col) - задаёт число бакетов и колонку(и) для хэширования. Spark вычисляет hash(col) % numBuckets для каждой строки и записывает её в соответствующий файл. Одинаковые значения col всегда попадают в один и тот же бакет.

sortBy(col) - дополнительная сортировка строк внутри каждого бакета по указанной колонке. Это не обязательно для устранения Shuffle, но убирает фазу Sort из Sort-Merge Join: данные уже отсортированы, Spark может сразу перейти к Merge. Без sortBy Shuffle всё равно устраняется, но Sort-фаза остаётся.

saveAsTable(name) - критически важный метод. Он не только записывает файлы, но и регистрирует метаданные в Hive Metastore: сколько бакетов, по какой колонке, с какой сортировкой. Без этой регистрации Spark не знает о бакетизации при планировании join.

Почему write.parquet(path) не работает

# ❌ НЕ РАБОТАЕТ: метаданные не сохраняются в Metastore
orders_df.write \
    .bucketBy(32, "user_id") \
    .sortBy("user_id") \
    .parquet("/data/orders_bucketed/")   # ← Parquet на диске есть, метаданных нет

# При последующем join:
spark.read.parquet("/data/orders_bucketed/").join(users_bucketed, "user_id")
# Spark видит просто Parquet без метаданных → выбирает SortMergeJoin с Shuffle

# ✅ РАБОТАЕТ: метаданные в Metastore
orders_df.write.bucketBy(32, "user_id").sortBy("user_id").saveAsTable("orders_bucketed")

# При чтении обязательно через spark.table()
orders = spark.table("orders_bucketed")   # ← Spark подтягивает метаданные из Metastore
users  = spark.table("users_bucketed")
result = orders.join(users, "user_id")    # ← Catalyst видит совместимые бакеты → нет Shuffle

Причина: write.parquet() записывает данные в файловую систему напрямую, минуя Metastore. Spark читает такой путь как обычный Parquet без каких-либо знаний о распределении строк по файлам. Только saveAsTable создаёт запись в Metastore с bucket-метаданными, и только spark.table() при чтении эти метаданные получает.

Чтение через spark.table() vs spark.read

# ✅ Правильно: Metastore вернёт bucket metadata → Plan содержит BucketedScan
orders = spark.table("orders_bucketed")
users  = spark.table("users_bucketed")

result = orders.join(users, "user_id")
result.explain(mode="simple")
# Physical Plan:
# SortMergeJoin [user_id], [user_id], Inner
# :- FileScan parquet default.orders_bucketed [...]
# :    SelectedBucketsCount: 32 out of 32   ← BucketedScan
# +- FileScan parquet default.users_bucketed [...]
#     SelectedBucketsCount: 32 out of 32    ← BucketedScan

# ❌ Неправильно: spark.read.table() и spark.read.parquet() не то же самое что spark.table()
orders_wrong = spark.read.format("parquet").load("/path/to/orders_bucketed")
# Метаданные не загружены → нет BucketedScan → будет Shuffle

5. Жёсткие условия срабатывания

Условие 1: Одинаковые bucket-ключи = join-ключи

Join в запросе должен выполняться строго по той же колонке, по которой таблицы забакетированы. Если join по другому ключу - бакеты не совпадают по смыслу, Shuffle неизбежен.

# ✅ Работает: join по user_id = bucket-ключ обеих таблиц
orders.join(users, "user_id")

# ❌ Не работает: join по order_id - orders забакетированы по user_id
orders.join(products, "product_id")
# orders: bucket по user_id, products: bucket по product_id
# Это разные ключи → совпадения нет → SortMergeJoin с Shuffle

# ❌ Не работает: join по выражению
orders.join(users, orders["user_id"] == F.abs(users["user_id"]))
# Spark не может свести F.abs(user_id) к bucket-ключу

Условие 2: Одинаковое число бакетов

Обе таблицы должны иметь ровно одинаковое число бакетов. Если числа не совпадают, Spark не может напрямую сопоставить бакеты - строки с одинаковым user_id окажутся в файлах с разными номерами в двух таблицах.

# ❌ Ломает Bucket Join: разное число бакетов
orders_df.write.bucketBy(32, "user_id").saveAsTable("orders_bucketed")
users_df.write.bucketBy(16, "user_id").saveAsTable("users_bucketed")  # 16 ≠ 32!

orders = spark.table("orders_bucketed")
users  = spark.table("users_bucketed")
result = orders.join(users, "user_id")
result.explain(mode="simple")
# Physical Plan:
# SortMergeJoin [user_id], [user_id], Inner
# :- Exchange hashpartitioning(user_id, 200)   ← SHUFFLE всё равно есть!
#    +- FileScan parquet default.orders_bucketed
# +- Exchange hashpartitioning(user_id, 200)   ← SHUFFLE есть!
#    +- FileScan parquet default.users_bucketed
# Spark не доверяет несовместимым бакетам и делает Shuffle поверх них

Что делать при разном числе бакетов?

Параметр spark.sql.bucketsInJoinReductionTolerate (Spark 3.x) позволяет Spark объединять бакеты если их число кратно. Например, если у одной таблицы 32 бакета, а у другой 64, Spark может объединить каждые 2 бакета большей таблицы в один и выполнить join без Shuffle. Условие: большая таблица должна иметь число бакетов, кратное числу бакетов меньшей.

# Включить join при кратном числе бакетов
spark.conf.set("spark.sql.bucketsInJoinReductionTolerate", "true")

# orders: 32 бакета, users: 64 бакета → 64/32=2 (кратно) → Spark объединит 2 бакета users
# Это работает, но добавляет накладные расходы на объединение

Условие 3: Одинаковый тип данных bucket-ключа

Тип данных join-ключа должен быть физически одинаковым в обеих таблицах. Неявное приведение типов (implicit cast) ломает Bucket Join - Spark выполняет cast и теряет возможность сопоставить бакеты.

# ❌ Ломает Bucket Join: user_id как Long в orders и Integer в users
orders_df.withColumn("user_id", F.col("user_id").cast("long")) \
         .write.bucketBy(32, "user_id").saveAsTable("orders_bucketed")

users_df.withColumn("user_id", F.col("user_id").cast("int")) \
        .write.bucketBy(32, "user_id").saveAsTable("users_bucketed")

# Spark при join увидит Long vs Integer → выполнит cast(int → long)
# hash(Long(100)) ≠ hash(Integer(100)) с Murmur3 в разных контекстах
# → бакеты не совпадают → Shuffle

# ✅ Правильно: одинаковый тип в обеих таблицах
orders_df.withColumn("user_id", F.col("user_id").cast("long")) \
         .write.bucketBy(32, "user_id").saveAsTable("orders_bucketed")

users_df.withColumn("user_id", F.col("user_id").cast("long")) \   # тот же тип!
        .write.bucketBy(32, "user_id").saveAsTable("users_bucketed")

Чтобы диагностировать проблему с типами - смотрите на DESCRIBE EXTENDED table и сравнивайте типы bucket-колонок в обеих таблицах.

Условие 4: Влияние sortBy на Sort-фазу SMJ

sortBy() при записи сортирует строки внутри каждого бакета. Это не обязательно для устранения Shuffle, но критично для устранения Sort-фазы в SMJ.

Сценарий Exchange (Shuffle) Sort (сортировка) Итоговый план
Нет бакетизации Есть Есть Exchange → Sort → Sort → Merge
bucketBy без sortBy Нет Есть ~~Exchange~~ → Sort → Sort → Merge
bucketBy + sortBy Нет Нет ~~Exchange~~ → ~~Sort~~ → Merge
# ─── Вариант A: только bucketBy, без sortBy ───────────────────────────────────
orders_df.write.bucketBy(32, "user_id").saveAsTable("orders_no_sort")
users_df.write.bucketBy(32, "user_id").saveAsTable("users_no_sort")

result_a = spark.table("orders_no_sort").join(spark.table("users_no_sort"), "user_id")
result_a.explain(mode="simple")
# SortMergeJoin [user_id], [user_id], Inner
# :- Sort [user_id ASC]           ← Sort остался!
#    +- FileScan (BucketedScan)   ← Shuffle убран
# +- Sort [user_id ASC]           ← Sort остался!
#    +- FileScan (BucketedScan)   ← Shuffle убран

# ─── Вариант B: bucketBy + sortBy ────────────────────────────────────────────
orders_df.write.bucketBy(32, "user_id").sortBy("user_id").saveAsTable("orders_sorted")
users_df.write.bucketBy(32, "user_id").sortBy("user_id").saveAsTable("users_sorted")

result_b = spark.table("orders_sorted").join(spark.table("users_sorted"), "user_id")
result_b.explain(mode="simple")
# SortMergeJoin [user_id], [user_id], Inner
# :- FileScan parquet orders_sorted (BucketedScan)   ← ни Shuffle, ни Sort!
# +- FileScan parquet users_sorted (BucketedScan)    ← ни Shuffle, ни Sort!

Совет: всегда используйте sortBy с тем же ключом что и bucketBy. Это не увеличивает время записи (сортировка выполняется в рамках той же shuffle-подобной операции при записи), но устраняет Sort при каждом join.

Условие 5: Совместимость при комбинации с партиционированием

Если таблицы одновременно партиционированы и забакетированы, обе должны иметь совместимые структуры партиций.

# ✅ Совместимые партиционирование + бакетизация
orders_df.write \
    .partitionBy("order_date") \        # одинаково
    .bucketBy(32, "user_id") \          # одинаково
    .sortBy("user_id") \                # одинаково
    .saveAsTable("orders_full")

users_df.write \
    .bucketBy(32, "user_id") \          # та же бакетизация
    .sortBy("user_id") \
    .saveAsTable("users_full")
# users не партиционированы по дате - это нормально
# При join Spark применит Partition Pruning к orders и Bucket Join к обеим

# ❌ Несовместимые партиционирования ломают Bucket Join
orders_df.write.partitionBy("order_date").bucketBy(32, "user_id").saveAsTable("orders_p")
users_df.write.partitionBy("user_country").bucketBy(32, "user_id").saveAsTable("users_p")
# Разные partition-ключи → Spark не может гарантировать co-location → Shuffle

6. Практика: читаем execution plan

Полный демо-пример

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

spark = SparkSession.builder \
    .appName("Bucketing-Practice") \
    .config("spark.sql.warehouse.dir", "/tmp/spark-warehouse-demo") \
    .config("spark.sql.autoBroadcastJoinThreshold", "-1")  # принудительно SMJ
    .enableHiveSupport() \
    .getOrCreate()

# ─── Генерация данных ─────────────────────────────────────────────────────────
N_ORDERS = 5_000_000
N_USERS  = 500_000

orders = (
    spark.range(N_ORDERS)
    .withColumn("user_id", (F.col("id") % N_USERS).cast("long"))
    .withColumn("amount",  (F.rand(seed=42) * 1000).cast("double"))
    .withColumn("status",  F.when(F.col("id") % 4 == 0, "cancelled").otherwise("paid"))
    .drop("id")
)

users = (
    spark.range(N_USERS)
    .withColumn("user_id",  F.col("id").cast("long"))
    .withColumn("country",  F.when(F.col("id") % 2 == 0, "RU").otherwise("US"))
    .withColumn("tier",     F.when(F.col("id") % 5 == 0, "vip").otherwise("regular"))
    .drop("id")
)

# ─── Запись: обе таблицы с одинаковой бакетизацией ───────────────────────────
print("Записываем бакетированные таблицы...")
orders.write.mode("overwrite") \
    .bucketBy(16, "user_id").sortBy("user_id") \
    .saveAsTable("orders_b16")

users.write.mode("overwrite") \
    .bucketBy(16, "user_id").sortBy("user_id") \
    .saveAsTable("users_b16")

# ─── Запись: без бакетизации (для сравнения) ─────────────────────────────────
orders.write.mode("overwrite").saveAsTable("orders_plain")
users.write.mode("overwrite").saveAsTable("users_plain")

print("Готово.")

Сравнение планов: с бакетизацией и без

# ─── ЗАПРОС: суммарные продажи по стране и тиру ──────────────────────────────
def build_query(orders_tbl, users_tbl):
    return (
        spark.table(orders_tbl)
        .filter(F.col("status") == "paid")
        .join(spark.table(users_tbl), "user_id")
        .groupBy("country", "tier")
        .agg(
            F.sum("amount").alias("total_amount"),
            F.count("*").alias("order_count"),
        )
    )

# ─── План без бакетизации ─────────────────────────────────────────────────────
print("=== ПЛАН БЕЗ БАКЕТИЗАЦИИ ===")
build_query("orders_plain", "users_plain").explain(mode="formatted")

# Вывод (ключевые части):
# == Physical Plan ==
# HashAggregate(...)
#  +- Exchange hashpartitioning(country, tier, 200)
#     +- HashAggregate(...)
#        +- SortMergeJoin [user_id], [user_id], Inner
#           :- Sort [user_id ASC NULLS FIRST]
#           :  +- Exchange hashpartitioning(user_id, 200)     ← SHUFFLE здесь!
#           :     +- Filter (status = paid)
#           :        +- FileScan parquet default.orders_plain
#           +- Sort [user_id ASC NULLS FIRST]
#              +- Exchange hashpartitioning(user_id, 200)     ← SHUFFLE здесь!
#                 +- FileScan parquet default.users_plain
#
# 3 узла Exchange = 3 shuffle операции

# ─── План с бакетизацией ─────────────────────────────────────────────────────
print("=== ПЛАН С БАКЕТИЗАЦИЕЙ ===")
build_query("orders_b16", "users_b16").explain(mode="formatted")

# Вывод (ключевые части):
# == Physical Plan ==
# HashAggregate(...)
#  +- Exchange hashpartitioning(country, tier, 200)    ← этот Shuffle - для groupBy
#     +- HashAggregate(...)
#        +- SortMergeJoin [user_id], [user_id], Inner
#           :- FileScan parquet default.orders_b16      ← НЕТ Exchange перед join!
#           :    PartitionFilters: []
#           :    SelectedBucketsCount: 16 out of 16    ← BucketedScan
#           +- FileScan parquet default.users_b16       ← НЕТ Exchange перед join!
#                SelectedBucketsCount: 16 out of 16    ← BucketedScan
#
# 1 узел Exchange = только для groupBy агрегации (неизбежен)
# Shuffle для join ПОЛНОСТЬЮ УСТРАНЁН

Ключевые маркеры в explain:

  • SelectedBucketsCount: N out of N - Spark читает бакетированные данные, BucketedScan активен.
  • Отсутствие Exchange hashpartitioning(user_id, ...) перед SortMergeJoin - Shuffle для join устранён.
  • Отсутствие Sort [user_id ASC] перед SortMergeJoin - данные уже отсортированы (sortBy при записи).

Диагностика «сломанного» бакетирования

# ─── ТЕСТ 1: разное число бакетов ────────────────────────────────────────────
users.write.mode("overwrite") \
    .bucketBy(8, "user_id").sortBy("user_id") \  # 8 вместо 16!
    .saveAsTable("users_b8")

build_query("orders_b16", "users_b8").explain(mode="simple")
# Результат: Exchange hashpartitioning появятся снова → Bucket Join не сработал

# ─── ТЕСТ 2: broadcast отключён, но AQE может вмешаться ─────────────────────
spark.conf.set("spark.sql.adaptive.enabled", "true")
spark.conf.set("spark.sql.autoBroadcastJoinThreshold", str(50 * 1024 * 1024))
# Если users_b16 маленькая (< 50 МБ) - AQE может переключить на BroadcastHashJoin
# В этом случае Shuffle тоже устраняется, но по другой причине
# Убедитесь что таблицы достаточно большие для честного теста

# ─── ТЕСТ 3: чтение через read.parquet() вместо spark.table() ────────────────
orders_wrong = spark.read.parquet("/tmp/spark-warehouse-demo/orders_b16")
users_wrong  = spark.read.parquet("/tmp/spark-warehouse-demo/users_b16")
orders_wrong.join(users_wrong, "user_id").explain(mode="simple")
# Exchange hashpartitioning появится! Метаданные не были загружены из Metastore.

Как убедиться что Bucket Join не "случайно" сломал AQE

# AQE может незаметно переключить SMJ → BHJ при маленьких данных
# Это маскирует отсутствие Bucket Join (потому что Shuffle тоже не видно)
# Для честного теста Bucket Join отключите AQE и Broadcast:

spark.conf.set("spark.sql.adaptive.enabled", "false")
spark.conf.set("spark.sql.autoBroadcastJoinThreshold", "-1")

# Теперь explain покажет честный план: есть ли Exchange или нет
build_query("orders_b16", "users_b16").explain(mode="simple")

7. Как рассчитать число бакетов

Выбор числа бакетов - один из самых практически важных вопросов. Нет универсальной формулы, но есть набор правил:

Правило 1: ориентир на объём данных

Каждый бакет должен содержать 128–256 МБ данных - такой же ориентир, как для shuffle-партиций. При меньшем размере бакета накладные расходы на открытие файла превышают полезную работу; при большем - не хватает памяти в Executor.

# Формула: numBuckets = ceil(tableSize / targetBucketSize)
table_size_gb = 500       # ГБ
target_bucket_mb = 256    # МБ = 0.25 ГБ
num_buckets = (table_size_gb * 1024) / target_bucket_mb  # ≈ 2000

# На практике округляем до ближайшей степени двойки:
# 2000 → 2048 (2^11)
# Использование степеней двойки упрощает работу с кратным числом бакетов

import math
def recommended_buckets(size_gb, target_mb=256):
    raw = (size_gb * 1024) / target_mb
    power = math.ceil(math.log2(raw))
    return 2 ** power

print(recommended_buckets(500))   # 2048
print(recommended_buckets(50))    # 256
print(recommended_buckets(10))    # 64

Правило 2: число должно быть одинаковым для всех join-таблиц

Если планируете join таблиц A (50 ГБ), B (500 ГБ) и C (5 ГБ), все три должны иметь одинаковое число бакетов. Число выбирается по самой большой таблице с округлением: 2048. A (50 ГБ) и C (5 ГБ) будут иметь файлы по 25 МБ и 2.5 МБ - это нормально, избыточно мелких файлов можно избежать через OPTIMIZE (Delta Lake) или compaction (Iceberg).

Правило 3: степени двойки и кратность

Когда таблицы имеют разный масштаб, удобно использовать кратное число бакетов: 64 / 128 / 256. Тогда Spark может объединять бакеты при join (при spark.sql.bucketsInJoinReductionTolerate = true), и несмотря на разное число бакетов, Shuffle может не понадобиться.


8. Взаимодействие с AQE

AQE (Adaptive Query Execution) и Bucketing иногда конфликтуют. AQE принимает решения в runtime, основываясь на реальных данных, и может переопределить bucket-план.

Когда AQE игнорирует Bucketing

AQE может решить объединить несколько маленьких shuffle-партиций (coalesce partitions) или переключить SMJ → BHJ. В обоих случаях исходный Bucket Join план меняется и нужно проверить через explain после выполнения (isFinalPlan=true).

spark.conf.set("spark.sql.adaptive.enabled", "true")

result = spark.table("orders_b16").join(spark.table("users_b16"), "user_id")
result.cache()
result.count()  # trigger execution

# Посмотреть ФИНАЛЬНЫЙ план после AQE-решений
result.explain(mode="formatted")
# AdaptiveSparkPlan isFinalPlan=true   ← это финальный план
# Проверьте: есть ли SelectedBucketsCount и нет ли Exchange для user_id join

spark.sql.sources.bucketing.autoBucketedScanEnabled

В Spark 3.1+ появился параметр, который позволяет AQE динамически отключать BucketedScan если он оказывается невыгодным (например, если нужны только несколько бакетов):

spark.conf.set("spark.sql.sources.bucketing.autoBucketedScanEnabled", "true")
# Если запрос фильтрует данные и нужно только 2 из 16 бакетов,
# Spark будет читать только эти 2 файла (Bucket Pruning)

9. Bucketing в Modern Lakehouse форматах

Delta Lake и Bucketing

Delta Lake использует собственные механизмы оптимизации, которые в большинстве случаев делают классический Hive Bucketing излишним:

Z-Order Clustering - физическая кластеризация файлов по часто используемым фильтровым колонкам. Spark читает только файлы, которые могут содержать нужные данные, основываясь на column statistics в Transaction Log. Этот механизм работает без фиксированного числа бакетов и не требует совпадения между таблицами.

Liquid Clustering (Delta Lake 3.1+) - адаптивная кластеризация, которая автоматически перераспределяет данные по файлам без жёсткого привязывания к числу бакетов. Это эволюция, которая устраняет главный недостаток классического bucketing - негибкость.

# Delta Lake Z-Order: аналог бакетизации для read-оптимизации
spark.sql("OPTIMIZE orders_delta ZORDER BY (user_id)")
# Данные перераспределяются по файлам чтобы похожие user_id были рядом
# → Data Skipping при фильтрации по user_id
# Но это НЕ даёт Shuffle-free join! Только помогает читать меньше данных.

# Delta Lake Liquid Clustering (3.1+): для write и join оптимизации
spark.sql("""
    CREATE TABLE orders_liquid
    USING DELTA
    CLUSTER BY (user_id)
    AS SELECT * FROM orders
""")

Apache Iceberg и Hidden Partitioning

Iceberg поддерживает скрытое партиционирование (PARTITIONED BY (bucket(N, col))), которое по сути является бакетизацией на уровне файловой организации. В отличие от Hive Bucketing, Iceberg-бакеты работают через Iceberg Catalog и не требуют Hive Metastore.

# Iceberg с bucket partitioning
spark.sql("""
    CREATE TABLE iceberg.orders (
        user_id BIGINT,
        amount  DOUBLE,
        status  STRING
    )
    USING iceberg
    PARTITIONED BY (bucket(32, user_id))
""")
# Строки с одинаковым user_id попадают в один бакет
# Iceberg при scan знает какие файлы нужны → читает только нужные бакеты
# НО: это не то же самое что Hive Bucketing - shuffle для join всё равно может быть

10. Anti-patterns

Anti-pattern 1: Over-bucketing - слишком много мелких файлов

# ❌ 1024 бакета для таблицы 500 МБ:
# 500 МБ / 1024 = 0.5 МБ на бакет - катастрофически мало!
small_df.write.bucketBy(1024, "user_id").saveAsTable("over_bucketed")
# Результат: 1024 крошечных файла по ~0.5 МБ
# NameNode/Metastore перегружен метаданными
# Открытие 1024 файлов медленнее чем один проход по 500 МБ
# Parquet footer overhead превышает полезные данные

# ✅ Правильно: ориентироваться на 128-256 МБ на бакет
# 500 МБ / 256 МБ = ~2 бакета → но минимум 4-8 для параллелизма
small_df.write.bucketBy(4, "user_id").saveAsTable("correctly_bucketed")

Anti-pattern 2: Бакетизация таблицы без Metastore (объектное хранилище)

Классический Hive Bucketing требует Hive Metastore. В чистых object-storage средах (S3, GCS без Hive Catalog) или при использовании spark.read.parquet() Bucket Join не работает. Если ваша инфраструктура не имеет Metastore - рассматривайте Delta Lake Liquid Clustering или Iceberg bucket partitioning как альтернативу.

Anti-pattern 3: Бакетизация таблиц с частым обновлением

# Проблема: бакеты не перераспределяются автоматически при append
new_orders.write.mode("append").bucketBy(32, "user_id").saveAsTable("orders_b32")
# Каждый append добавляет новые файлы с суффиксами _00000..._00031
# Через год: тысячи файлов per bucket → медленный open, медленный read
# Нужна регулярная компакция (OPTIMIZE в Delta) для поддержки здоровья бакетов

# Признаки деградации: один бакет = много мелких файлов
spark.sql("DESCRIBE DETAIL orders_b32")
# numFiles резко выросло → пора компактить

Anti-pattern 4: Случайный выбор ключа бакетизации

# ❌ Бакетизация по колонке с плохим распределением (skew)
# Если 80% заказов от топ-1000 пользователей:
orders.write.bucketBy(32, "user_id").saveAsTable("orders_skewed")
# Бакеты с популярными user_id будут в 100 раз больше остальных
# SMJ будет медленным из-за дисбаланса → Spill на некоторых Executor-ах

# ✅ Проверьте распределение ключа перед бакетизацией:
orders.groupBy("user_id").count().orderBy(F.col("count").desc()).show(20)
# Если топ-N ключей содержат > 30% данных - скорее всего skew проблема
# Рассмотрите salt technique или другой ключ бакетизации

11. Best Practices и чек-лист

Когда Bucketing действительно полезен

Bucketing - это low-level оптимизация. Её применение оправдано только в конкретных сценариях:

  • Обе таблицы большие (> 10–50 ГБ) - слишком большие для Broadcast Join
  • Join по одному и тому же ключу выполняется регулярно (десятки или сотни раз в день)
  • Таблицы стабильны - схема и ключи не меняются месяцами
  • Используется plain Parquet без продвинутого Lakehouse формата
  • Приоритет - долгоживущие аналитические витрины, а не потоковые пайплайны

Когда Bucketing НЕ нужен

  • Одна из таблиц < 100–500 МБ - Broadcast Join проще и не требует Metastore
  • AQE + хорошая статистика справляется с оптимизацией join автоматически
  • Использование Delta Lake или Iceberg - у них есть более гибкие аналоги
  • Streaming-пайплайны - Structured Streaming не поддерживает saveAsTable с bucketBy
  • Таблицы часто обновляются или имеют эволюцию схемы

Чек-лист перед применением Bucketing

Проектирование:

  • Определить join-ключ с хорошим распределением (проверить на skew)
  • Одинаковое число бакетов во всех join-таблицах (степень двойки)
  • Размер бакета ≈ 128–256 МБ
  • Используйте sortBy с тем же ключом что и bucketBy

Запись:

  • Обязательно saveAsTable(), не write.parquet()
  • Собрать статистику: ANALYZE TABLE ... COMPUTE STATISTICS FOR ALL COLUMNS
  • Проверить результат: DESCRIBE EXTENDED table_name → секция Num Buckets

Верификация:

  • Отключить AQE и Broadcast для чистого теста: autoBroadcastJoinThreshold = -1
  • Запустить explain(mode="formatted") и найти SelectedBucketsCount + отсутствие Exchange
  • Сравнить время выполнения и Shuffle Read в Spark UI до/после бакетизации

Обслуживание:

  • Настроить регулярную компакцию при append-нагрузке
  • Мониторить число файлов per bucket (DESCRIBE DETAIL)

Домашнее задание

Дано:

# Две таблицы для join
events     = spark.read.parquet("s3a://data/events/")      # 200 ГБ, event_user_id, event_type, ts
user_attrs = spark.read.parquet("s3a://data/user_attrs/")  # 5 ГБ, user_id, segment, region

Текущий запрос выполняется 30 минут, Shuffle Read = 190 ГБ:

result = events.join(user_attrs, events["event_user_id"] == user_attrs["user_id"]) \
               .groupBy("segment", "region") \
               .agg(F.count("*").alias("event_count"), F.countDistinct("event_user_id").alias("users"))
result.write.parquet("s3a://output/segment_stats/")

Задание 1 - Диагностика. Запустите explain(mode="formatted") и найдите узлы Exchange hashpartitioning. Сколько их? Какой из них соответствует join, а какой - groupBy?

Задание 2 - Применить Bucketing. Перепишите pipeline: сохраните обе таблицы как забакетированные Hive-таблицы. Рассчитайте число бакетов на основе размеров таблиц (целевой размер бакета - 256 МБ).

Задание 3 - Убедиться в результате. Запустите новый запрос через spark.table() с отключённым AQE и Broadcast. Найдите SelectedBucketsCount в explain и подтвердите отсутствие Exchange перед join.

Задание 4 - Ответить на вопрос. Почему ключ join events["event_user_id"] == user_attrs["user_id"] может помешать Bucket Join? Как переписать join чтобы Spark мог использовать бакеты? Подсказка: посмотрите на имена колонок.


В следующем уроке разберём Storage Partitioning: как организовать файлы на диске для максимального Partition Pruning и минимального overhead при чтении в Parquet, Delta Lake и Iceberg.