Broadcast Join: autoBroadcastJoinThreshold, hints и Task Too Large

Полный разбор Broadcast Join в Spark: анатомия BroadcastHashJoin vs SortMergeJoin, механизм broadcast переменных, autoBroadcastJoinThreshold и проблемы Size Estimation, явные Hints в PySpark и SQL, ошибки OOM на Driver и Executor, Task Too Large, ANALYZE TABLE, влияние AQE и мониторинг через Spark UI.

optimization

1. Анатомия Broadcast Join: почему Join — самая дорогая операция

В аналитических запросах JOIN почти всегда является главным источником затрат производительности. Чтобы понять почему и как Broadcast Join решает эту проблему, нужно разобрать физику стандартного Join'а.

Sort-Merge Join: почему Shuffle так дорого стоит

В стандартном SortMergeJoin (SMJ) Spark выполняет следующие шаги:

Схема показывает ключевую проблему SMJ: обе таблицы перераспределяются по сети. Если таблица фактов содержит 100 GB и таблица измерений — 2 GB, то по сети передаётся суммарно 102 GB. При медленной сети (10 Gbps) только передача займёт 80+ секунд.

Плюс операция Sort для обеих сторон. Итого SMJ для больших таблиц — это:

  • 2× Shuffle Write (обе стороны пишут отсортированные данные на диск)
  • 2× Shuffle Read (обе стороны читают с других Executor'ов)
  • 2× Sort (сортировка каждой стороны)
  • 1× Merge (слияние потоков)

Broadcast Hash Join (BHJ) устраняет Shuffle полностью для большой стороны:

Почему BHJ быстрее: большая таблица (100 GB) вообще не перемещается. Каждый Executor читает её локально со своего DataNode (NODE_LOCAL) и сразу делает hash lookup в локальной копии малой таблицы. Сетевой трафик = только рассылка малой таблицы (2 GB × число Executor'ов).

Сравнительная стоимость для типичного production запроса

Операция SMJ (100 GB fact + 2 GB dim) BHJ (100 GB fact + 2 GB dim)
Shuffle Write 102 GB 0 GB (large table не shuffleится)
Shuffle Read 102 GB 0 GB
Sort 2 × O(N log N) 0 (нет сортировки)
Broadcast 2 GB × N Executors
Join Merge O(N) Hash Probe O(N)
Итоговое ускорение baseline 5–50×

2. Механизм Broadcast Variables: как работает рассылка

Понимание механизма broadcast переменных помогает правильно настраивать параметры и диагностировать проблемы.

Жизненный цикл Broadcast переменной

Диаграмма показывает несколько важных деталей:

Фаза 1 (Collect): малая таблица полностью собирается в память Driver'а через collect(). Это первое место возможного OOM: если таблица 10 GB, Driver должен иметь 10+ GB свободной памяти.

Фаза 3 (Distribute): Spark использует BitTorrent-подобный протокол распределения. Executor'ы не все качают с Driver'а — они передают данные друг другу. Это снижает нагрузку на Driver при большом числе Executor'ов.

Фаза 4 (Deserialize): каждый Executor десериализует байты и строит Hash Table. Это второе место возможного OOM: каждому Executor'у нужно дополнительно broadcast_size в памяти для Hash Table.

Сколько памяти потребляет Broadcast

from pyspark.sql import SparkSession, functions as F

spark = SparkSession.builder.getOrCreate()

# Оценка потребления памяти при broadcast
def estimate_broadcast_memory(df_size_mb: float, n_executors: int) -> dict:
    """
    Оценивает суммарное потребление памяти для broadcast.

    При broadcast каждый Executor хранит полную копию таблицы
    + Driver держит оригинал для рассылки.

    Сериализованный размер обычно меньше in-memory благодаря
    сжатию (Snappy/LZ4 для Parquet данных).
    """
    # Предполагаем 2x overhead: in-memory > on-disk
    in_memory_size_mb = df_size_mb * 2

    # Driver: оригинальные данные + сериализованная копия
    driver_memory_mb = in_memory_size_mb + df_size_mb

    # Каждый Executor: сериализованные данные + Hash Table (2x overhead)
    executor_memory_per_node_mb = in_memory_size_mb * 2

    # Суммарное потребление на кластере
    total_cluster_mb = driver_memory_mb + executor_memory_per_node_mb * n_executors

    return {
        "table_size_mb": df_size_mb,
        "driver_memory_mb": driver_memory_mb,
        "per_executor_memory_mb": executor_memory_per_node_mb,
        "total_cluster_memory_mb": total_cluster_mb,
        "recommendation": (
            "SAFE" if df_size_mb < 100
            else "CAUTION (проверьте executor.memory)" if df_size_mb < 500
            else "DANGEROUS (рассмотрите SMJ)"
        )
    }

# Пример: broadcast 50 MB таблицы на кластер из 20 Executor'ов
result = estimate_broadcast_memory(50, 20)
for k, v in result.items():
    print(f"  {k}: {v}")
# Driver: ~150 MB overhead
# Per Executor: ~100 MB
# Total cluster: 150 + 100×20 = 2150 MB = ~2 GB дополнительной памяти

3. autoBroadcastJoinThreshold: автоматический выбор стратегии

Catalyst Optimizer автоматически выбирает BHJ вместо SMJ, если оценочный размер одной стороны JOIN меньше порога autoBroadcastJoinThreshold.

Как работает автоматический выбор

# По умолчанию: 10 MB (10 * 1024 * 1024 = 10485760 байт)
spark.conf.get("spark.sql.autoBroadcastJoinThreshold")
# '10485760'

# Проверяем текущий порог:
threshold_bytes = int(spark.conf.get("spark.sql.autoBroadcastJoinThreshold"))
print(f"Auto-broadcast порог: {threshold_bytes / 1024 / 1024:.0f} MB")

# Изменить порог:
spark.conf.set("spark.sql.autoBroadcastJoinThreshold", str(100 * 1024 * 1024))  # 100 MB

# Или при создании SparkSession:
spark = SparkSession.builder \
    .config("spark.sql.autoBroadcastJoinThreshold", "200MB") \
    .getOrCreate()

# Полностью отключить auto-broadcast (SMJ для всех JOIN):
spark.conf.set("spark.sql.autoBroadcastJoinThreshold", "-1")

Что использует Catalyst для оценки размера:

  • Parquet/Delta/Iceberg файлы: читает статистику из file footers (размер файлов на диске × estimated compression ratio)
  • Hive таблицы: читает totalSize из HMS (если запускали ANALYZE TABLE)
  • После трансформаций: применяет коэффициенты из Cost-Based Optimizer (CBO) — но часто ошибается!

Проблема Size Estimation: когда Catalyst теряет размер

Это главная причина почему auto-broadcast не срабатывает даже для очевидно маленьких таблиц.

# Сценарий: таблица 2 GB фильтруется до 5 MB
# Catalyst не знает точный размер после фильтрации!

dim_full = spark.read.parquet("hdfs://cluster/dims/regions/")  # 2 GB

# Фильтр убирает 99.75% данных → остаётся 5 MB
dim_filtered = dim_full.filter(F.col("region") == "EU")

# Catalyst оценивает: размер до фильтра × selectivity
# Если нет статистики по колонке region → selectivity = дефолтный (например 0.1-0.33)
# Оценка: 2 GB × 0.1 = 200 MB >> 10 MB threshold → НЕТ auto-broadcast!
# Реальный размер: 5 MB << 10 MB threshold → ДОЛЖЕН быть broadcast

# Проверяем что думает Catalyst:
dim_filtered.explain("cost")
# EstimatedSizes(sizeInBytes=200.0 MiB ← НЕПРАВИЛЬНАЯ оценка!)

Типичные сценарии где Size Estimation ошибается:

  • Фильтр после WITH CTE или подзапроса
  • Сложные условия с несколькими OR/AND
  • Результат DISTINCT или LIMIT
  • Данные из UDF (Catalyst не знает что вернёт UDF)
  • Результат агрегации без CBO

4. Явные Hints: принудительный Broadcast

Когда автоматика подводит — используем явные подсказки (Hints). Они принудительно меняют физический план, игнорируя autoBroadcastJoinThreshold.

Синтаксис Hints в PySpark

from pyspark.sql import SparkSession, functions as F

spark = SparkSession.builder \
    .config("spark.sql.adaptive.enabled", "true") \
    .getOrCreate()

# Данные
orders = spark.table("bronze.orders")    # большая таблица: 500 GB
regions = spark.table("silver.regions")  # малая таблица: 5 MB

# ── Способ 1: F.broadcast() — наиболее распространённый ────────────────
# Оборачиваем малую сторону в broadcast()
result = orders.join(
    F.broadcast(regions),   # Hint: broadcast regions
    on="region_id",
    how="inner"
)

# ── Способ 2: DataFrame.hint() ─────────────────────────────────────────
# Более явный синтаксис, особенно полезен при читабельности кода
result = orders.join(
    regions.hint("broadcast"),   # Hint непосредственно на DataFrame
    on="region_id",
    how="inner"
)

# ── Способ 3: hint() с именованными таблицами ──────────────────────────
# Полезно когда нет прямого доступа к DataFrame объекту
result = spark.sql("""
    SELECT o.*, r.region_name
    FROM orders o
    INNER JOIN regions r
    ON o.region_id = r.id
""").hint("broadcast", "r")  # Hint через имя таблицы/алиаса

# Проверяем что BHJ выбран:
result.explain("formatted")
# Expected output:
# == Physical Plan ==
# *BroadcastHashJoin [region_id#1], [id#15], Inner, BuildRight, false
# :- *Project [...]
# :  +- *Scan parquet bronze.orders
# +- BroadcastExchange HashedRelationBroadcastMode(...)
#    +- *Project [...]
#       +- *Scan parquet silver.regions

Синтаксис SQL Hints

-- Spark SQL hints (рекомендуемый синтаксис)
SELECT /*+ BROADCAST(r) */
    o.order_id,
    o.amount,
    r.region_name
FROM orders o
JOIN regions r ON o.region_id = r.id;

-- Альтернативные синтаксисы (все работают в Spark SQL):
SELECT /*+ MAPJOIN(r) */ ...    -- старый Hive синтаксис
SELECT /*+ BROADCASTJOIN(r) */ ... -- ещё один вариант

-- Для нескольких таблиц:
SELECT /*+ BROADCAST(d1), BROADCAST(d2) */
    f.value,
    d1.name,
    d2.category
FROM facts f
JOIN dim1 d1 ON f.dim1_id = d1.id
JOIN dim2 d2 ON f.dim2_id = d2.id;

Антипаттерн: broadcast через collect() в замыкании

# ❌ НЕПРАВИЛЬНО: передаём данные через замыкание Python
# Это не broadcast в понимании Spark — это TaskSerialization!
lookup_data = spark.table("silver.regions").collect()  # → Python список
lookup_dict = {row.id: row.region_name for row in lookup_data}

@udf(returnType=StringType())
def get_region_name(region_id):
    return lookup_dict.get(region_id, "unknown")  # замыкание!

# Проблема: lookup_dict (~5 MB) сериализуется в КАЖДУЮ Task
# При 1000 Task'ах: 5 GB только на передачу словаря
# Ошибка: "Serialized task exceeds 100 MiB limit"

# ✅ ПРАВИЛЬНО вариант А: Spark broadcast через F.broadcast()
# Полностью избегаем UDF и используем нативный JOIN
result = facts.join(
    F.broadcast(spark.table("silver.regions")),
    "region_id"
)

# ✅ ПРАВИЛЬНО вариант Б: SparkContext.broadcast() для Python UDF
# Если UDF всё же необходим
sc = spark.sparkContext
broadcast_dict = sc.broadcast(lookup_dict)  # Spark управляет рассылкой

@udf(returnType=StringType())
def get_region_name_safe(region_id):
    return broadcast_dict.value.get(region_id, "unknown")  # broadcast.value!

# broadcast.value читается локально — нет TaskSerialization!

5. Риски OOM на Driver и Executor

Broadcast Join при неправильном применении превращается из оптимизации в источник катастрофических падений.

OOM на Driver: первая волна

Driver собирает всю малую таблицу в памяти перед рассылкой. Если таблица на самом деле не такая маленькая:

# Сценарий: инженер думает что таблица 50 MB, но она 2 GB
# (данные выросли за полгода, а hint остался в коде)

large_dim = spark.table("dims.user_segments")  # ВЫРОСЛА до 2 GB!

# Ошибка в логах Driver'а:
# ERROR SparkContext: Error while broadcasting
# java.lang.OutOfMemoryError: Java heap space
#     at java.util.Arrays.copyOf(Arrays.java:3210)
#     at org.apache.spark.sql.execution.joins.HashedRelation$.apply(...)

# Или более загадочно:
# ERROR TransportResponseHandler: Still have 200 requests outstanding
#   when connection from /executor3 is closed
# → Executor не смог получить broadcast из-за проблем с Driver

# Диагностика: проверяем реальный размер
def check_dataframe_size(df, sample_fraction=0.01):
    """Оценивает реальный размер DataFrame через выборку."""
    sample_count = df.sample(fraction=sample_fraction, seed=42).count()
    full_count = df.count()
    # Получаем размер через explain статистику
    df.createOrReplaceTempView("_size_check")
    stats = spark.sql("ANALYZE TABLE _size_check COMPUTE STATISTICS")
    return full_count

# Предохранитель от OOM при ручном broadcast:
SAFE_BROADCAST_LIMIT_MB = 500

def safe_broadcast(df, name: str = "unnamed"):
    """
    Безопасный broadcast с проверкой размера.
    Выбрасывает исключение если таблица слишком большая.
    """
    # Оцениваем через Catalyst
    catalyst_size = df._jdf.queryExecution().optimizedPlan().stats().sizeInBytes()
    size_mb = catalyst_size / 1024 / 1024

    if size_mb > SAFE_BROADCAST_LIMIT_MB:
        raise ValueError(
            f"DataFrame '{name}' слишком большой для broadcast: {size_mb:.0f} MB. "
            f"Лимит: {SAFE_BROADCAST_LIMIT_MB} MB. "
            f"Рассмотрите SMJ или увеличьте spark.driver.memory."
        )

    print(f"Broadcast '{name}': {size_mb:.1f} MB — безопасно")
    return F.broadcast(df)

OOM на Executor: вторая волна

Даже если Driver справился с рассылкой, каждый Executor должен удержать Hash Table в памяти параллельно с обработкой данных:

Executor Memory Layout при BHJ:
- Execution Memory: обработка строк большой таблицы
- Storage Memory: ← BHJ Hash Table живёт ЗДЕСЬ! (broadcast переменная)
- User Memory: UDF, метаданные

Если Hash Table (5 GB) > Storage Memory (5.91 GB):
→ Hash Table не помещается полностью
→ Partial eviction на диск
→ OOM или очень медленный join
# Типичный лог OOM на Executor при большом broadcast:
WARN  TransportRequestHandler: Error while invoking RPC for broadcast
ERROR Executor: Failed to fetch broadcast 42
  SparkException: Size of serialized broadcast value > spark.driver.maxResultSize
  (2147483648 bytes). Consider using a smaller dataset or increasing
  spark.driver.maxResultSize.

# Или при нехватке памяти на Executor:
ERROR Container: Container killed by YARN for exceeding memory limits.
  20.4 GB of 20 GB physical memory used.
  Container killed for exceeding memory limits!

Правила безопасного broadcast:

# Практические лимиты (не жёсткие, зависят от конфигурации):
LIMITS = {
    "safe":       100,   # MB — всегда безопасно
    "caution":    500,   # MB — проверьте executor.memory
    "risky":     1000,   # MB — нужно явно настраивать параметры
    "dangerous": 2000,   # MB — скорее всего SMJ лучше
}

# Если хотите broadcast таблицу 500 MB на кластере с 20 GB executors:
# 1. Storage Memory per Executor ≈ 5.91 GB (см. Memory Fractions урок)
# 2. Hash Table для 500 MB таблицы ≈ 1-2 GB в RAM
# 3. 1-2 GB / 5.91 GB = 17-34% Storage Memory — приемлемо

# Но при 50 executor'ах × 20 GB memory × 500 MB broadcast:
# Суммарное потребление: 50 × 2 GB = 100 GB дополнительной памяти на кластере

6. Ошибка Task Too Large: TaskSerialization

Task Too Large — специфическая ошибка, возникающая когда Task Description превышает лимит сериализации. Это отличается от OOM при broadcast.

Причины Task Too Large

# ❌ Антипаттерн 1: большой объект в замыкании UDF
large_model_weights = {...}  # ML модель: 200 MB Python dict

@udf(returnType=FloatType())
def predict(features):
    return large_model_weights["weights"].dot(features)  # closure!

# Spark включает large_model_weights в каждую Task Description
# При 200 Tasks: 200 × 200 MB = 40 GB только на Task сериализацию
# Ошибка: Serialized task exceeds max allowed size (100 MB)

# ❌ Антипаттерн 2: collect() результата в map()
lookup = spark.table("dim").collect()  # 500K строк → Python list
df.filter(df.id.isin([row.id for row in lookup]))  # isin с огромным списком!
# Python list с 500K элементов сериализуется в каждую Task

# Признаки проблемы в логах:
# WARN TaskSetManager: Stage 3 contains a task of very large size (1234 KiB).
#   The maximum recommended task size is 100 KiB.
# ERROR SparkContext: Failed to get broadcast variable 42
#   java.lang.RuntimeException: Serialized task exceeds max allowed size

Правильные решения для Task Too Large

# ✅ Решение 1: Spark native JOIN вместо collect() + isin()
valid_ids = spark.table("dim").select("id")  # малый DataFrame
result = df.join(F.broadcast(valid_ids), "id", "semi")
# Semi JOIN: сохраняет строки df где id есть в valid_ids (как isin)
# Никакого collect(), никаких замыканий

# ✅ Решение 2: SparkContext.broadcast() для Python UDF
# Данные рассылаются через Spark механизм, а не через Task serialization
sc = spark.sparkContext

large_model = {
    "weights": [0.1, 0.2, 0.3, 0.4, 0.5],
    "bias": 0.05,
    "threshold": 0.5
}
broadcast_model = sc.broadcast(large_model)

@udf(returnType=FloatType())
def predict_safe(features):
    model = broadcast_model.value  # читаем через broadcast, не closure!
    weights = model["weights"]
    return float(sum(w * f for w, f in zip(weights, features)))

# ✅ Решение 3: Pandas UDF с broadcast через closure-безопасный паттерн
# Для ML моделей — загружаем модель ОДИН РАЗ на Executor, а не в каждую Task
from pyspark.sql.functions import pandas_udf
import pandas as pd

model_path = "/shared/models/lgbm_model.pkl"  # путь, а не объект

@pandas_udf("float")
def predict_lgbm(features: pd.Series) -> pd.Series:
    import pickle
    # Модель загружается один раз при первом вызове на Executor
    # (lazy loading через Python functools.lru_cache или модульная переменная)
    import importlib.util
    return pd.Series([0.5] * len(features))  # placeholder

7. ANALYZE TABLE: помогаем Catalyst принимать правильные решения

Самый надёжный способ обеспечить автоматическое применение broadcast — поддерживать актуальную статистику таблиц.

Как ANALYZE TABLE помогает

-- Сбор статистики для таблицы (размер, количество строк, nullCount)
ANALYZE TABLE silver.regions COMPUTE STATISTICS;

-- Сбор статистики по конкретным колонкам (гистограммы распределений)
-- Это позволяет CBO точнее оценивать selectivity фильтров
ANALYZE TABLE silver.regions COMPUTE STATISTICS
    FOR COLUMNS region_id, country_code, region_name;

-- Для партиционированных таблиц:
ANALYZE TABLE bronze.orders PARTITION (dt='2024-01-15')
    COMPUTE STATISTICS;

-- Проверяем что статистика собрана:
DESCRIBE EXTENDED silver.regions;
-- Statistics: 2048 bytes, 50 rows (теперь Catalyst знает точный размер)
# В PySpark: сбор статистики программно
spark.sql("ANALYZE TABLE silver.regions COMPUTE STATISTICS")

# После ANALYZE TABLE:
dim = spark.table("silver.regions")
dim.explain("cost")
# Теперь EstimatedSizes(sizeInBytes=2.0 KiB) — ПРАВИЛЬНАЯ оценка!
# Catalyst корректно выбирает BHJ вместо SMJ

# Автоматизация через Airflow: запускать ANALYZE после каждой загрузки dim
def analyze_dimension_tables(spark, dim_tables: list[str]) -> None:
    """
    Обновляет статистику таблиц-измерений после ETL.
    Вызывается как финальный шаг DAG загрузки dimensions.
    """
    for table in dim_tables:
        print(f"ANALYZE TABLE {table}...")
        spark.sql(f"ANALYZE TABLE {table} COMPUTE STATISTICS")
        # Для Delta Lake / Iceberg можно добавить:
        # spark.sql(f"ANALYZE TABLE {table} COMPUTE STATISTICS FOR ALL COLUMNS")
        print(f"  Done!")

analyze_dimension_tables(spark, [
    "dims.regions",
    "dims.products",
    "dims.customers",
])

Когда ANALYZE TABLE не помогает

ANALYZE TABLE актуален только для физических таблиц. Для промежуточных DataFrame (результат фильтрации, агрегации) нужно или использовать Hints, или включить CBO:

spark = SparkSession.builder \
    # Cost-Based Optimizer: использует статистику колонок для join reordering
    .config("spark.sql.cbo.enabled", "true") \

    # Join Reorder: CBO переупорядочивает JOIN'ы для минимальной стоимости
    .config("spark.sql.cbo.joinReorder.enabled", "true") \

    # Размер таблицы при котором CBO собирает статистику автоматически
    .config("spark.sql.statistics.autoUpdate.enabled", "true") \

    .getOrCreate()

8. AQE и Broadcast Join: динамическое переключение

Adaptive Query Execution (AQE) в Spark 3.x кардинально меняет работу с Broadcast Join. Теперь Spark может переключиться на BHJ в процессе выполнения, даже если изначально запланировал SMJ.

Как AQE динамически выбирает BHJ

# AQE Broadcast Join конфигурация:
spark = SparkSession.builder \
    .config("spark.sql.adaptive.enabled", "true") \

    # Порог для динамического upgrade SMJ → BHJ
    # AQE использует этот порог ПОСЛЕ выполнения предыдущего Stage
    # (реальные размеры, а не оценки!)
    .config("spark.sql.adaptive.autoBroadcastJoinThreshold", "30MB") \

    # Порог для статического выбора BHJ (до выполнения)
    .config("spark.sql.autoBroadcastJoinThreshold", "10MB") \

    .getOrCreate()

# Демонстрация AQE upgrade:
regions = spark.table("silver.regions")  # 800 MB (Catalyst думает)

# После фильтра — фактически 3 MB
regions_eu = regions.filter(F.col("continent") == "EU")

# БЕЗ AQE: SMJ (800 MB > 10 MB threshold → нет broadcast)
# С AQE: Stage 1 завершился → AQE видит 3 MB → upgrade to BHJ!

orders = spark.table("bronze.orders")
result = orders.join(regions_eu, "region_id")

# Проверяем что AQE применил BHJ:
result.explain("formatted")
# После выполнения:
# AdaptiveSparkPlan isFinalPlan=true
# +- BroadcastHashJoin (replaced SortMergeJoin after runtime statistics)

9. Мониторинг через Spark UI: верификация Broadcast

SQL/DataFrame вкладка: чтение Physical Plan

В Spark UI → SQL/DataFrame вкладка нужно найти правильные операторы:

Успешный BHJ:

BroadcastHashJoin [region_id#1], [id#15], Inner, BuildRight
:- *Project [order_id#0, amount#2, region_id#1]
:  +- *FileScan parquet bronze.orders
+- BroadcastExchange HashedRelationBroadcastMode(...)
   +- *Filter (isnotnull(id#15) AND (...)
      +- *FileScan parquet silver.regions

Метрика BroadcastExchange:
  time to broadcast: 0.8 sec
  data size: 2.3 MB

Неудачный SMJ когда должен быть BHJ:

SortMergeJoin [region_id#1], [id#15], Inner   ← НЕТ broadcast!
:- Sort [region_id#1 ASC]
:  +- Exchange hashpartitioning(region_id#1, 200)   ← Shuffle!
:     +- Project [...]
:        +- FileScan parquet bronze.orders
+- Sort [id#15 ASC]
   +- Exchange hashpartitioning(id#15, 200)   ← Shuffle для dim!
      +- Filter (...)
         +- FileScan parquet silver.regions

Программная проверка типа JOIN

def get_join_strategy(df) -> str:
    """
    Определяет тип JOIN из Physical Plan текстового вывода.
    Полезно для автоматической верификации в тестах.
    """
    plan = df._jdf.queryExecution().executedPlan().toString()

    if "BroadcastHashJoin" in plan:
        return "BroadcastHashJoin"
    elif "SortMergeJoin" in plan:
        return "SortMergeJoin"
    elif "BroadcastNestedLoopJoin" in plan:
        return "BroadcastNestedLoopJoin"
    elif "CartesianProduct" in plan:
        return "CartesianProduct (осторожно!)"
    else:
        return "Unknown"

# Использование:
result = orders.join(F.broadcast(regions), "region_id")
strategy = get_join_strategy(result)
print(f"Join strategy: {strategy}")

# В тестах:
assert strategy == "BroadcastHashJoin", \
    f"Ожидался BHJ но получен {strategy}. Проверьте размер regions!"

Ключевые метрики в Spark UI для BHJ

BroadcastExchange метрики:

  • time to broadcast — сколько времени заняла рассылка. Если > 5 сек → таблица слишком большая или сеть перегружена
  • data size — реальный размер переданных данных. Должен быть меньше autoBroadcastJoinThreshold
  • number of output rows — число строк в broadcast таблице

Сравнение времени BHJ vs SMJ:

BHJ: 0.8 сек broadcast + 2.1 сек join = 2.9 сек total
SMJ: 45 сек shuffle write + 38 сек shuffle read + 12 сек sort + 8 сек merge = 103 сек

BHJ выигрыш: 103 / 2.9 = 35x ускорение!


10. Best Practices и антипаттерны

Полный набор правильных практик

from pyspark.sql import SparkSession, functions as F

spark = SparkSession.builder \
    # Автоматический broadcast для таблиц < 100 MB
    .config("spark.sql.autoBroadcastJoinThreshold", "100MB") \

    # AQE для динамического upgrade
    .config("spark.sql.adaptive.enabled", "true") \
    .config("spark.sql.adaptive.autoBroadcastJoinThreshold", "200MB") \

    # CBO для лучших оценок
    .config("spark.sql.cbo.enabled", "true") \

    .getOrCreate()

# ── Паттерн 1: Fact + Dimension ────────────────────────────────────────────
# Классический случай: справочник всегда маленький
def join_with_dimension(
    fact_df,
    dim_table: str,
    join_key: str,
    use_broadcast: bool = True,
    max_dim_size_mb: float = 500.0,
) -> "DataFrame":
    """
    Безопасный join фактовой таблицы с таблицей-измерением.

    Автоматически применяет broadcast если dim < max_dim_size_mb.
    """
    dim_df = spark.table(dim_table)

    # Проверяем размер dim (через Catalyst stats)
    dim_size_bytes = dim_df._jdf.queryExecution().optimizedPlan() \
        .stats().sizeInBytes()
    dim_size_mb = dim_size_bytes / 1024 / 1024

    if use_broadcast and dim_size_mb < max_dim_size_mb:
        print(f"BHJ: {dim_table} = {dim_size_mb:.1f} MB (< {max_dim_size_mb} MB)")
        return fact_df.join(F.broadcast(dim_df), join_key)
    else:
        print(f"SMJ: {dim_table} = {dim_size_mb:.1f} MB (≥ {max_dim_size_mb} MB)")
        return fact_df.join(dim_df, join_key)


# ── Паттерн 2: Broadcast Bloom Filter (Spark 3.x) ─────────────────────────
# Для очень больших dim таблиц (нельзя broadcast полностью):
# Bloom Filter передаётся вместо полной таблицы для early filtering
spark.conf.set("spark.sql.optimizer.runtime.bloomFilter.enabled", "true")
spark.conf.set("spark.sql.optimizer.runtime.bloomFilter.creationSideThreshold", "10MB")
# При 10 MB bloom filter → отфильтровывает большинство не-совпадений
# ещё до полного SMJ → значительно меньше данных в Shuffle

# ── Паттерн 3: Reuse broadcast переменной ─────────────────────────────────
# Один и тот же dim используется в нескольких JOIN'ах
regions = spark.table("silver.regions").cache()  # кешируем в Storage Memory
regions.count()  # материализуем кеш

# Теперь при broadcast — данные читаются из Storage Memory, не с диска
result1 = fact1.join(F.broadcast(regions), "region_id")
result2 = fact2.join(F.broadcast(regions), "region_id")  # переиспользуем кеш
result3 = fact3.join(F.broadcast(regions), "region_id")

Матрица выбора стратегии JOIN

Размер малой таблицы Ситуация Рекомендация
< 10 MB Авто-broadcast работает Ничего не делать
10-100 MB Авто-broadcast не всегда config("autoBroadcastJoinThreshold", "100MB")
100-500 MB Авто не работает, ручной риск F.broadcast() + проверить executor.memory
> 500 MB Broadcast опасен SMJ или Bloom Filter Join
Неизвестно, результат трансформаций Catalyst ошибается ANALYZE TABLE + AQE

Антипаттерны, которые нужно запомнить

# ❌ АНТИПАТТЕРН 1: broadcast растущей таблицы фактов
# Таблица фактов с 1B+ строк НИКОГДА не должна быть broadcast
facts = spark.table("bronze.orders")  # 500 GB!
dims = spark.table("silver.regions")  # 5 MB

# НЕПРАВИЛЬНО: broadcast большой таблицы
result = F.broadcast(facts).join(dims, "region_id")  # OOM!

# ПРАВИЛЬНО: broadcast маленькой таблицы
result = facts.join(F.broadcast(dims), "region_id")  # ✓

# ❌ АНТИПАТТЕРН 2: broadcast без проверки размера в prod коде
# В prod данные могут вырасти! Добавляйте проверки:
def safe_broadcast_join(fact_df, dim_df, key, max_mb=500):
    dim_size_mb = dim_df._jdf.queryExecution().optimizedPlan() \
        .stats().sizeInBytes() / 1024 / 1024
    if dim_size_mb > max_mb:
        raise RuntimeError(f"dim_df слишком большой: {dim_size_mb:.0f} MB > {max_mb} MB")
    return fact_df.join(F.broadcast(dim_df), key)

# ❌ АНТИПАТТЕРН 3: отключать broadcast полностью
# spark.conf.set("spark.sql.autoBroadcastJoinThreshold", "-1")
# → все JOIN'ы становятся SMJ → катастрофическое замедление

# ❌ АНТИПАТТЕРН 4: broadcast в Structured Streaming
# Broadcast join плохо работает со стримингом —
# broadcast таблица не обновляется при изменении данных
# Используйте Stream-Static JOIN (без broadcast) или foreachBatch

Лабораторное задание: диагностика и оптимизация JOIN

# lab_broadcast_join.py
from pyspark.sql import SparkSession, functions as F
import time

spark = SparkSession.builder \
    .master("local[8]") \
    .appName("broadcast-join-lab") \
    .config("spark.sql.adaptive.enabled", "false") \  # выключаем AQE для демонстрации
    .getOrCreate()

spark.sparkContext.setLogLevel("WARN")

# Создаём тестовые данные
n_facts = 5_000_000
n_dims = 1_000  # маленький справочник

facts = spark.range(n_facts).select(
    F.col("id").alias("order_id"),
    (F.rand() * n_dims).cast("int").alias("region_id"),
    (F.rand() * 1000).alias("amount")
)

dims = spark.range(n_dims).select(
    F.col("id").alias("region_id"),
    F.concat(F.lit("Region_"), F.col("id")).alias("region_name")
)

# Кешируем исходные данные
facts.cache()
dims.cache()
facts.count()
dims.count()

print("=== ТЕСТ 1: SortMergeJoin (без broadcast) ===")
spark.conf.set("spark.sql.autoBroadcastJoinThreshold", "-1")  # отключаем broadcast
t0 = time.time()
result_smj = facts.join(dims, "region_id") \
    .groupBy("region_name").sum("amount")
count_smj = result_smj.count()
t_smj = time.time() - t0
print(f"SMJ: {t_smj:.1f}с ({count_smj} строк)")

print("\n=== ТЕСТ 2: BroadcastHashJoin (ручной hint) ===")
spark.conf.set("spark.sql.autoBroadcastJoinThreshold", "-1")  # всё равно выключен
t0 = time.time()
result_bhj = facts.join(F.broadcast(dims), "region_id") \
    .groupBy("region_name").sum("amount")
count_bhj = result_bhj.count()
t_bhj = time.time() - t0
print(f"BHJ: {t_bhj:.1f}с ({count_bhj} строк)")

print(f"\n=== РЕЗУЛЬТАТЫ ===")
print(f"SMJ: {t_smj:.1f}с")
print(f"BHJ: {t_bhj:.1f}с")
print(f"Ускорение: {t_smj/t_bhj:.1f}x")
print()
print("Стратегии JOIN (проверьте .explain()):")
print(f"SMJ strategy: {get_join_strategy(facts.join(dims, 'region_id'))}")
print(f"BHJ strategy: {get_join_strategy(facts.join(F.broadcast(dims), 'region_id'))}")

spark.stop()

Итоги: золотые правила Broadcast Join

Правило 1: Broadcast работает только для малой таблицы. Всегда broadcast'ите dimension/lookup, никогда — fact таблицу.

Правило 2: Проверяйте реальный размер, а не оценочный. Catalyst часто ошибается для отфильтрованных датасетов. ANALYZE TABLE + AQE дают точные оценки.

Правило 3: Проверяйте план через .explain(). Убедитесь что в Physical Plan есть BroadcastHashJoin, а не SortMergeJoin.

Правило 4: AQE решает большинство проблем автоматически. В Spark 3.2+ с включённым AQE многие ситуации где раньше нужны были Hints — теперь решаются автоматически.

Правило 5: Лимит безопасности. 200-500 MB — разумный верхний предел для broadcast при типичных конфигурациях. Выше — проверяйте executor.memory и время broadcast.

Правило 6: Никогда не используйте collect() + isin() вместо JOIN. Это антипаттерн с Task Serialization. Используйте Semi JOIN с broadcast вместо.