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