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()