Cartesian Join Trap: случайный cross join, BNLJ на TB и non-equi join ловушки

Случайный Cartesian join - верный способ создать O(N×M) строк и подвесить кластер на сутки. Разбираем три пути к неожиданному декартовому произведению и как их обнаружить через EXPLAIN.

optimization core

Как возникает случайный Cartesian Join

Путь 1: JOIN без условия

# ❌ Забыли условие join - получили Cartesian Product
result = orders.join(customers)
# Нет "on" параметра → crossJoin → N × M строк
# orders: 100M строк, customers: 1M строк → 100 ТРИЛЛИОНОВ строк

# Spark по умолчанию ЗАПРЕЩАЕТ это (если spark.sql.crossJoin.enabled = false)
# Но с включённым crossJoin.enabled - молча выполняется
# ✅ Всегда указывайте условие
result = orders.join(customers, orders.customer_id == customers.id)

Путь 2: JOIN с общим именем колонки - неожиданный cross join в SQL

-- ❌ Таблицы не имеют общих колонок → CROSS JOIN!
SELECT * FROM orders JOIN customers
-- Без ON и без USING → Cartesian в ANSI SQL

-- ✅ Обязательно ON или USING
SELECT * FROM orders JOIN customers ON orders.customer_id = customers.id

Путь 3: Non-equi JOIN → Broadcast Nested Loop Join (BNLJ)

# ❌ Non-equi условие → Spark выбирает BNLJ (похоже на Cartesian)
df1.join(df2, df1.start_date <= df2.end_date)
# Сложность O(N × M) - для 10M × 10M строк = 100 ТРИЛЛИОНОВ операций сравнения
# Диагностика: EXPLAIN покажет BroadcastNestedLoopJoin
df1.join(df2, df1.start_date <= df2.end_date).explain()
# == Physical Plan ==
# BroadcastNestedLoopJoin BuildRight, Inner, (start_date <= end_date)
#   ...

Как обнаружить до запуска

# Всегда делайте explain() для новых join'ов
df.explain(mode="formatted")

# Ключевые слова в плане = немедленно разобраться:
# CartesianProduct          ← явный cross join
# BroadcastNestedLoopJoin   ← non-equi join или join без условий
# Exchange(repl)            ← shuffle_replicate_nl hint
-- В Spark SQL:
EXPLAIN FORMATTED
SELECT * FROM a JOIN b ON a.dt <= b.dt

-- Ищем: BroadcastNestedLoopJoin
-- Нормально только если ОБЕ таблицы маленькие (< 10 MB каждая)

Безопасные альтернативы для non-equi join

Диапазонный join через Bucketing + Range

# Вместо df1.join(df2, (df1.ts >= df2.start) & (df1.ts < df2.end))
# → Разбить на range buckets + equi-join по bucket

from pyspark.sql.functions import (date_trunc, col)

# Округлить ts до часа → equi-join по часу + range-фильтр внутри
df1 = df1.withColumn("hour", date_trunc("hour", "ts"))
df2 = df2.withColumn("hour", date_trunc("hour", "start"))

df1.join(df2, "hour").filter(
    (col("ts") >= col("start")) & (col("ts") < col("end"))
)
# Shuffle только по hour (equi), range-фильтр локальный

Interval join через Iceberg или специализированные библиотеки

# Для очень больших таблиц с interval join - рассмотреть:
# 1. Предфильтрацию: оставить только overlapping периоды
# 2. Broadcast если одна сторона мала
# 3. Разбить задачу на range buckets

# Broadcast маленькой таблицы при non-equi
small_df = spark.read.parquet("ranges/")   # < 10 MB
big_df.join(broadcast(small_df),
    (big_df.ts >= small_df.start) & (big_df.ts < small_df.end))
# BNLJ, но с broadcast → без shuffle, быстро при малой таблице

Когда Cartesian/BNLJ допустим

# ✅ Допустимо: оба датасета маленькие и контролируемые
calendar = spark.range(365).toDF("day")   # 365 строк
products = spark.createDataFrame([(1,), (2,), (3,)], ["product_id"])

# Генерация grid: 365 × 3 = 1095 строк - нормально
calendar.crossJoin(products)

# ✅ Допустимо: broadcast non-equi join, правая сторона < 10 MB
events.join(broadcast(small_ranges),
    (events.ts >= small_ranges.start) & (events.ts < small_ranges.end))

Защита от случайного Cartesian

# Оставить crossJoin.enabled = false (дефолт) - Spark выбросит исключение
spark.conf.set("spark.sql.crossJoin.enabled", "false")  # default

# Тогда любой crossJoin потребует явного вызова:
df1.crossJoin(df2)   # явно - значит намеренно

# ✅ Добавить это в CI/CD проверку explain():
def has_cartesian(plan_str):
    return "CartesianProduct" in plan_str or "BroadcastNestedLoopJoin" in plan_str

plan = df.explain(extended=True, mode="string")
assert not has_cartesian(plan), "Accidental Cartesian join detected!"

Диагностика после запуска

Spark UI → Stages:
  Input: 10 GB
  Output: 50 TB    ← очевидный взрыв данных
  Task Duration: running for 6+ hours

Stages → SQL Plan:
  CartesianProduct или BroadcastNestedLoopJoin + большой Input Size