Bucketing: sort-merge join без shuffle - настройка и ограничения
Как бакетизация переносит стоимость shuffle с каждого запроса на момент записи данных, физический механизм Bucket Co-location в Sort-Merge Join, жёсткие условия срабатывания и архитектурные ограничения в современных Lakehouse.
В предыдущем уроке мы разобрали 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:
- Открывает только директорию
event_date=2026-01-01/- Partition Pruning. - Внутри неё читает только нужные бакеты - Bucket Co-location.
- 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.