explode() misuse: N×M взрыв строк и когда использовать flatten

explode() превращает один массив в несколько строк. При неосторожном применении - особенно вложенном - количество строк взрывается экспоненциально.

core optimization

Что делает explode()

explode() разворачивает массив или map в отдельные строки - по одной на элемент:

from pyspark.sql.functions import explode, col

df = spark.createDataFrame([
    (1, ["a", "b", "c"]),
    (2, ["x", "y"]),
], ["id", "tags"])

df.withColumn("tag", explode("tags")).show()
# id | tag
# 1  | a
# 1  | b
# 1  | c
# 2  | x
# 2  | y
# Из 2 строк → 5 строк

Паттерн-убийца 1: двойной explode

# ❌ Вложенный explode на двух массивах
df = spark.createDataFrame([
    (1, ["a","b","c"], ["x","y","z"]),
], ["id", "col1", "col2"])

# col1 содержит 3 элемента, col2 - 3 элемента
# 1 строка → 3 строки (первый explode)
# 3 строки → 9 строк (второй explode) - это CARTESIAN!
df.withColumn("v1", explode("col1")) \
  .withColumn("v2", explode("col2")) \
  .count()  # = 9, а не 3

# На реальных данных: 1000 elements × 1000 elements = 1M строк из одной!
# ✅ Правильно - arrays_zip перед explode
from pyspark.sql.functions import arrays_zip

df.withColumn("zipped", arrays_zip("col1", "col2")) \
  .withColumn("pair", explode("zipped")) \
  .select("id", col("pair.col1").alias("v1"), col("pair.col2").alias("v2"))
# 1 строка → 3 строки (без картезиана)

Паттерн-убийца 2: explode → groupBy без необходимости

# ❌ explode только чтобы потом посчитать размер массива
df.withColumn("tag", explode("tags")) \
  .groupBy("id") \
  .count()

# ✅ Используйте size() - нет лишнего shuffle
from pyspark.sql.functions import size
df.withColumn("tag_count", size("tags"))

Паттерн-убийца 3: explode на nested arrays без flatten

# ❌ array of arrays - explode даёт массив, не элементы
df = spark.createDataFrame([
    (1, [["a","b"], ["c"]]),
], ["id", "nested"])

# Первый explode → [["a","b"], ["c"]] → ["a","b"] и ["c"]
# Второй explode → отдельные строки → двойной проход
df.withColumn("inner", explode("nested")) \
  .withColumn("elem", explode("inner"))

# ✅ flatten → один explode
from pyspark.sql.functions import flatten
df.withColumn("flat", flatten("nested")) \
  .withColumn("elem", explode("flat"))

Когда explode уместен

# ✅ Нормализация: раскрыть массив тегов для поиска
logs.withColumn("tag", explode("tags")) \
    .groupBy("tag").count()  # частота тегов

# ✅ JSON-массив в строки для последующего join
events.withColumn("item", explode(from_json("items", item_schema))) \
      .join(catalog, col("item.product_id") == catalog.id)

# ✅ posexplode - если нужен индекс элемента
from pyspark.sql.functions import posexplode
df.select("id", posexplode("tags").alias("pos", "tag"))

Альтернативы explode для агрегаций над массивами

from pyspark.sql.functions import aggregate, transform, exists, forall

# Сумма элементов массива без explode
df.withColumn("total", aggregate("amounts", lit(0), lambda acc, x: acc + x))

# Фильтр элементов массива без explode
df.withColumn("big", filter("amounts", lambda x: x > 100))

# Есть ли элемент больше 1000?
df.withColumn("has_big", exists("amounts", lambda x: x > 1000))

# Трансформация каждого элемента без explode
df.withColumn("doubled", transform("amounts", lambda x: x * 2))

Диагностика

# Посмотреть, сколько строк возвращает explode:
df.withColumn("elem", explode("array_col")).count()

# Если count() после explode >> исходного count() в разы —
# проверить среднее и максимальное количество элементов в массиве:
from pyspark.sql.functions import avg, max as spark_max, size
df.select(
    avg(size("array_col")).alias("avg_len"),
    spark_max(size("array_col")).alias("max_len")
).show()