LATERAL VIEW и explode: unnesting массивов в строки

Разворачиваем вложенные структуры: explode-семейство, LATERAL VIEW, GenerateExec, карательный Cartesian и higher-order functions как альтернатива.

core

Проблема вложенных структур

В реальных данных вложенность - норма, а не исключение. Event-потоки из Kafka хранят список тегов в одном поле. Заказ из интернет-магазина содержит массив позиций. Пользовательский профиль может хранить список адресов или историю действий. Всё это прекрасно ложится в Parquet как array<T> или array<struct<...>>, позволяя избежать нормализации на этапе сырых данных.

Но аналитика требует другого: нужно посчитать выручку по каждому SKU, найти самые популярные теги, построить воронку по шагам сессии. Для этого вложенные строки нужно "развернуть" - превратить одну строку с массивом из 5 элементов в 5 строк. Именно это делает семейство explode.

Одна строка превращается в три - по числу элементов в массиве. order_id при этом дублируется в каждую строку. Это фундаментальная операция нормализации на лету, без изменения Bronze-хранилища.

Почему не хранить сразу в нормализованном виде

Хранить массивы в одном поле выгодно по нескольким причинам:

  • Parquet columnar encoding: массив хранится компактно как одна колонка с repetition levels, без дублирования order_id на диске
  • Атомарность записи: весь заказ пишется в одной транзакции, не нужно координировать несколько таблиц
  • Schema evolution: добавить новое поле в struct проще, чем мигрировать нормализованную схему
  • Streaming latency: Kafka event - единица передачи, разбивать его нецелесообразно

Нормализация происходит при чтении (в Silver-слое), а не при записи. Это классический паттерн Medallion Architecture.


Функции семейства explode

В Spark есть несколько функций для разворачивания массивов. Каждая решает свою задачу.

explode()

Базовая функция. Принимает ArrayType или MapType, возвращает по одной строке на элемент. Строки с null в массиве удаляются.

from pyspark.sql import SparkSession
from pyspark.sql.functions import explode, col

spark = SparkSession.builder.appName("explode-demo").getOrCreate()

data = [
    (1, ["python", "spark", "kafka"]),
    (2, ["scala", "java"]),
    (3, None),          # null массив - строка удалится
    (4, []),            # пустой массив - строка удалится
]

df = spark.createDataFrame(data, ["user_id", "skills"])

df.select("user_id", explode("skills").alias("skill")).show()
+-------+------+
|user_id| skill|
+-------+------+
|      1|python|
|      1| spark|
|      1| kafka|
|      2| scala|
|      2|  java|
+-------+------+

Пользователи 3 и 4 исчезли из результата. Это поведение по умолчанию - аналог INNER JOIN. Если нужно сохранить строки с пустыми/null-массивами, используйте explode_outer.

Для MapType функция возвращает две колонки - ключ и значение:

from pyspark.sql.functions import explode
from pyspark.sql.types import MapType, StringType, IntegerType

data = [(1, {"apples": 3, "bananas": 5}), (2, {"oranges": 2})]
df = spark.createDataFrame(data, ["id", "basket"])

df.select("id", explode("basket").alias("fruit", "count")).show()
+---+-------+-----+
| id|  fruit|count|
+---+-------+-----+
|  1| apples|    3|
|  1|bananas|    5|
|  2|oranges|    2|
+---+-------+-----+

explode_outer()

Идентична explode(), но строки с null и пустыми массивами сохраняются - значение колонки становится null. Это аналог LEFT JOIN.

from pyspark.sql.functions import explode_outer

df.select("user_id", explode_outer("skills").alias("skill")).show()
+-------+------+
|user_id| skill|
+-------+------+
|      1|python|
|      1| spark|
|      1| kafka|
|      2| scala|
|      2|  java|
|      3|  null|   ← null массив сохранён
|      4|  null|   ← пустой массив сохранён
+-------+------+

Используйте explode_outer когда:

  • Нужно сохранить записи без вложенных данных для статистики
  • Строите LEFT JOIN-подобную семантику
  • Считаете пользователей без навыков как отдельную категорию

posexplode()

Добавляет к результату колонку с индексом элемента в массиве (0-based). Удобно когда порядок элементов имеет значение - шаги воронки, история действий, приоритет.

from pyspark.sql.functions import posexplode

data = [(1, ["impression", "click", "purchase"])]
df = spark.createDataFrame(data, ["session_id", "funnel_steps"])

df.select("session_id", posexplode("funnel_steps").alias("step_num", "step")).show()
+----------+--------+----------+
|session_id|step_num|      step|
+----------+--------+----------+
|         1|       0|impression|
|         1|       1|     click|
|         1|       2|  purchase|
+----------+--------+----------+

Индекс позволяет анализировать последовательность: найти, на каком шаге пользователи чаще всего уходят, или вычислить время между шагами если добавить timestamp-массив.

posexplode_outer()

Комбинация: сохраняет строки с null/пустыми массивами и добавляет индекс. При null-массиве оба поля (pos и col) будут null.

from pyspark.sql.functions import posexplode_outer

data = [(1, ["a", "b"]), (2, None)]
df = spark.createDataFrame(data, ["id", "vals"])

df.select("id", posexplode_outer("vals").alias("pos", "val")).show()
+---+----+----+
| id| pos| val|
+---+----+----+
|  1|   0|   a|
|  1|   1|   b|
|  2|null|null|
+---+----+----+

inline() - разворачивание массива структур

inline() - специализированная функция для array<struct<...>>. В отличие от explode, она сразу разворачивает поля структуры в отдельные колонки, без необходимости делать дополнительный select с col("item.field").

from pyspark.sql.functions import inline
from pyspark.sql.types import *

schema = StructType([
    StructField("order_id", IntegerType()),
    StructField("items", ArrayType(StructType([
        StructField("sku", StringType()),
        StructField("qty", IntegerType()),
        StructField("price", DoubleType()),
    ])))
])

data = [(1001, [("SKU-A", 2, 99.9), ("SKU-B", 1, 249.0)]),
        (1002, [("SKU-C", 3, 49.5)])]

df = spark.createDataFrame(data, schema)

df.select("order_id", inline("items")).show()
+--------+-----+---+-----+
|order_id|  sku|qty|price|
+--------+-----+---+-----+
|    1001|SKU-A|  2| 99.9|
|    1001|SKU-B|  1|249.0|
|    1002|SKU-C|  3| 49.5|
+--------+-----+---+-----+

Заметьте: поля структуры (sku, qty, price) стали отдельными колонками автоматически. С explode пришлось бы писать:

from pyspark.sql.functions import explode, col

df.select(
    "order_id",
    explode("items").alias("item")
).select(
    "order_id",
    col("item.sku"),
    col("item.qty"),
    col("item.price")
).show()

inline экономит один шаг и чище выражает намерение. inline_outer() - версия с сохранением строк без данных.

Сводная таблица функций

Функция NULL/пустой массив Возвращает позицию Разворачивает struct
explode удаляет строку нет нет
explode_outer сохраняет (null) нет нет
posexplode удаляет строку да нет
posexplode_outer сохраняет (null) да нет
inline удаляет строку нет да
inline_outer сохраняет (null) нет да

LATERAL VIEW в Spark SQL

LATERAL VIEW - это SQL-синтаксис для применения генерирующей функции к каждой строке таблицы. Функционально эквивалентен explode() в DataFrame API, но иногда удобнее при работе с несколькими вложенными массивами в одном запросе.

Базовый синтаксис

SELECT
    o.order_id,
    item.sku,
    item.qty,
    item.price
FROM orders o
LATERAL VIEW explode(o.items) tmp AS item

Здесь:

  • LATERAL VIEW говорит Spark что нужно применить функцию explode к каждой строке orders
  • tmp - псевдоним таблицы (virtual table alias), обязателен по синтаксису
  • AS item - имя колонки с результатом

tmp - это просто требование синтаксиса; фактически эта таблица нигде не используется. Некоторые диалекты SQL опускают её, но Spark требует явного псевдонима.

LATERAL VIEW с posexplode

SELECT
    session_id,
    step_pos,
    step_name
FROM sessions
LATERAL VIEW posexplode(funnel_steps) tmp AS step_pos, step_name
WHERE step_name = 'purchase'

Несколько LATERAL VIEW

Можно применить LATERAL VIEW несколько раз к разным колонкам в одном запросе:

SELECT
    user_id,
    skill,
    lang
FROM users
LATERAL VIEW explode(skills) t1 AS skill
LATERAL VIEW explode(languages) t2 AS lang

Каждый LATERAL VIEW применяется последовательно. Важно: если один массив имеет 3 элемента, а другой - 2, результат - декартово произведение (3 × 2 = 6 строк на исходную строку). Подробнее об этой ловушке - в разделе про Double Explode.

LATERAL VIEW OUTER

Аналог explode_outer в SQL-синтаксисе:

SELECT
    user_id,
    skill
FROM users
LATERAL VIEW OUTER explode(skills) tmp AS skill

Строки с null или пустыми массивами сохраняются, skill будет null.


Физический план: оператор GenerateExec

Чтобы понять что происходит внутри Spark при explode, посмотрим на физический план.

from pyspark.sql.functions import explode

df_exploded = df.select("order_id", explode("items").alias("item"))
df_exploded.explain(mode="formatted")

Упрощённый вывод:

== Physical Plan ==
Generate explode(items), [order_id], false, [item]
+- Scan parquet orders

Ключевой оператор - Generate (в старом формате) или GenerateExec в физическом плане. Это специализированный оператор Spark, который обрабатывает генерирующие функции.

Параметры GenerateExec

В плане вы увидите несколько атрибутов оператора:

  • generator: сама функция - explode, posexplode, inline и т.д.
  • requiredChildOutput: колонки из родительского датафрейма, которые нужно пробросить (например, order_id)
  • outer: true если используется explode_outer / LATERAL VIEW OUTER
  • generatorOutput: имена выходных колонок генератора

Как GenerateExec работает

GenerateExec работает как flatMap на уровне строк:

  1. Читает строку из дочернего оператора
  2. Вычисляет генерирующую функцию для заданной колонки
  3. Для каждого элемента результата - выдаёт одну выходную строку
  4. Пробрасывает requiredChildOutput в каждую выходную строку

Это означает: оператор не требует shuffle. GenerateExec - локальная операция, выполняется на каждом partition независимо. Это хорошая новость с точки зрения производительности.

Плохая новость: если массивы сильно неравномерны по размеру, partition skew возникнет после explode - одни партиции получат намного больше строк, чем другие.


Взрыв кардинальности (Cardinality Explosion)

explode умножает количество строк. Это не баг - это цель операции. Но важно понимать масштаб умножения.

Пример: в таблице events 1 миллиард строк. Среднее количество тегов на событие - 5, максимальное - 200. После explode(tags):

  • Среднее: 5 млрд строк
  • Максимум на партицию: зависит от того, как распределены "жирные" строки
# Посчитаем распределение размеров массивов до explode
from pyspark.sql.functions import size, percentile_approx, max as spark_max, avg

df.select(
    avg(size("items")).alias("avg_items"),
    spark_max(size("items")).alias("max_items"),
    percentile_approx(size("items"), 0.95).alias("p95_items"),
    percentile_approx(size("items"), 0.99).alias("p99_items"),
).show()

Это важный первый шаг перед разворачиванием больших таблиц: понять, во сколько раз вырастет датасет и не будет ли некоторые партиции несбалансированы.

Partition skew после explode

Если размер массивов сильно варьируется (некоторые пользователи имеют 1 тег, а другие - 10 000), после explode возникнет partition skew. Партиции с "жирными" строками будут содержать в тысячи раз больше данных.

# Перед repartition - оцениваем skew
from pyspark.sql.functions import explode, spark_partition_id, count

df_exploded = df.select("user_id", explode("tags").alias("tag"))

# Сколько строк в каждой партиции после explode?
df_exploded.groupBy(spark_partition_id()).count().orderBy("count", ascending=False).show(20)

Признаки проблемы: одна партиция в 10+ раз больше остальных. Решение:

# Перераспределить после explode
df_normalized = df.select("user_id", explode("tags").alias("tag"))
df_normalized = df_normalized.repartition(200)  # или по ключу

Или использовать salting если нужна группировка по user_id после explode.


Double Explode: антипаттерн

Самая опасная ловушка при работе с explode - применить его дважды к разным массивам одной строки без промежуточной агрегации.

Что происходит

from pyspark.sql.functions import explode

data = [(1, ["python", "spark"], ["backend", "data"])]
df = spark.createDataFrame(data, ["user_id", "skills", "roles"])

# ОПАСНО: двойной explode
df.select(
    "user_id",
    explode("skills").alias("skill"),
    explode("roles").alias("role")
).show()
+-------+------+-------+
|user_id| skill|   role|
+-------+------+-------+
|      1|python|backend|
|      1|python|   data|
|      1| spark|backend|
|      1| spark|   data|
+-------+------+-------+

Пользователь 1 имел 2 навыка и 2 роли. Результат - 4 строки (2 × 2). Это декартово произведение массивов. Для массивов длиной 100 и 50 - 5000 строк вместо ожидаемых 100 или 50.

Когда Double Explode это ошибка

Если массивы не связаны позиционно (это просто два независимых списка атрибутов), декартово произведение - логическая ошибка. Например:

  • Пользователь с тегами ["python", "spark"] и городами ["Moscow", "SPb"] - 4 комбинации не имеют смысла

Как избежать: взрывайте по одному

Правило: разворачивайте только один массив за раз. Если нужны оба массива - сделайте два отдельных разворачивания и объедините по ключу:

skills_df = df.select("user_id", explode("skills").alias("skill"))
roles_df  = df.select("user_id", explode("roles").alias("role"))

# Если нужна нормализованная таблица фактов по навыкам и роли
# JOIN только если есть смысловая связь между позициями
result = skills_df.join(roles_df, "user_id")  # только если это уместно

Когда Double Explode корректен

Единственный случай, когда два explode дают правильный результат - когда вы намеренно хотите декартово произведение. Например: все возможные комбинации категорий и регионов для таблицы "покрытия". Но это редкий случай, и его нужно документировать явно.

Альтернатива: posexplode + JOIN по позиции

Если два массива связаны позиционно (элемент 0 первого массива соответствует элементу 0 второго), используйте posexplode + JOIN по позиции:

from pyspark.sql.functions import posexplode

# Предположим: timestamps[i] соответствует events[i]
ts_df    = df.select("session_id", posexplode("timestamps").alias("pos", "ts"))
event_df = df.select("session_id", posexplode("events").alias("pos", "event"))

result = ts_df.join(event_df, ["session_id", "pos"])

Теперь элементы соединяются по индексу, а не декартово.


Практический кейс: e-commerce заказы

Разберём реалистичный сценарий: Bronze-таблица с заказами хранит позиции как array<struct>. Нужно построить Silver-таблицу с нормализованными позициями.

Структура данных

from pyspark.sql.types import *

order_schema = StructType([
    StructField("order_id", LongType()),
    StructField("customer_id", LongType()),
    StructField("order_date", StringType()),
    StructField("status", StringType()),
    StructField("items", ArrayType(StructType([
        StructField("sku", StringType()),
        StructField("product_name", StringType()),
        StructField("category", StringType()),
        StructField("qty", IntegerType()),
        StructField("unit_price", DoubleType()),
        StructField("discount_pct", DoubleType()),
    ]))),
])

data = [
    (1001, 42, "2024-01-15", "completed", [
        ("SKU-001", "Laptop Pro 15", "electronics", 1, 89999.0, 0.05),
        ("SKU-002", "Mouse Wireless", "electronics", 2, 1299.0, 0.0),
        ("SKU-003", "Laptop Bag", "accessories", 1, 3499.0, 0.1),
    ]),
    (1002, 17, "2024-01-15", "completed", [
        ("SKU-004", "Python Book", "books", 3, 599.0, 0.15),
    ]),
    (1003, 99, "2024-01-15", "pending", None),  # заказ без позиций
]

orders_df = spark.createDataFrame(data, order_schema)

Bronze → Silver: нормализация позиций

from pyspark.sql.functions import inline_outer, col, round as spark_round, to_date

# inline_outer сохраняет заказ 1003 (без позиций)
silver_items = orders_df.select(
    col("order_id"),
    col("customer_id"),
    to_date("order_date").alias("order_date"),
    col("status"),
    inline_outer("items"),  # разворачивает и сразу выводит поля структуры
)

# Добавим вычисляемые поля
silver_items = silver_items.withColumn(
    "line_total",
    spark_round(col("qty") * col("unit_price") * (1 - col("discount_pct")), 2)
)

silver_items.show(truncate=False)
+--------+-----------+----------+---------+-------+--------------+-----------+---+----------+------------+----------+
|order_id|customer_id|order_date|   status|    sku|  product_name|   category|qty|unit_price|discount_pct|line_total|
+--------+-----------+----------+---------+-------+--------------+-----------+---+----------+------------+----------+
|    1001|         42|2024-01-15|completed|SKU-001| Laptop Pro 15|electronics|  1|  89999.0|        0.05|  85499.05|
|    1001|         42|2024-01-15|completed|SKU-002|Mouse Wireless|electronics|  2|   1299.0|         0.0|   2598.00|
|    1001|         42|2024-01-15|completed|SKU-003|    Laptop Bag|accessories|  1|   3499.0|         0.1|   3149.10|
|    1002|         17|2024-01-15|completed|SKU-004|   Python Book|      books|  3|    599.0|        0.15|   1527.45|
|    1003|         99|2024-01-15|  pending|   null|          null|       null|  0|      null|        null|      null|
+--------+-----------+----------+---------+-------+--------------+-----------+---+----------+------------+----------+

Заказ 1003 сохранён со null в полях позиции - благодаря inline_outer.

Аналитика поверх Silver

Теперь стандартные запросы работают без вложенности:

from pyspark.sql.functions import sum as spark_sum, countDistinct

# Выручка по категориям
silver_items.filter(col("status") == "completed") \
    .groupBy("category") \
    .agg(
        spark_sum("line_total").alias("revenue"),
        countDistinct("order_id").alias("orders"),
        spark_sum("qty").alias("units_sold")
    ) \
    .orderBy("revenue", ascending=False) \
    .show()
+-----------+---------+------+----------+
|   category|  revenue|orders|units_sold|
+-----------+---------+------+----------+
|electronics|91596.05|     1|         3|
|accessories| 3149.10|     1|         1|
|      books| 1527.45|     1|         3|
+-----------+---------+------+----------+

Higher-order functions как альтернатива explode

Иногда explode - избыточная операция. Если цель - трансформация или фильтрация массива (но не нормализация), лучше использовать higher-order functions: они работают внутри массива, не увеличивая количество строк.

Spark предоставляет четыре основных higher-order function для массивов:

transform()

Применяет лямбду к каждому элементу массива, возвращает новый массив той же длины.

from pyspark.sql.functions import transform, col

data = [(1, [100.0, 200.0, 350.0])]
df = spark.createDataFrame(data, ["order_id", "prices"])

# Применить скидку 10% ко всем ценам
df.select(
    "order_id",
    transform("prices", lambda x: x * 0.9).alias("discounted_prices")
).show()
+--------+-----------------------+
|order_id|      discounted_prices|
+--------+-----------------------+
|       1|[90.0, 180.0, 315.0]  |
+--------+-----------------------+

Никакого умножения строк - один заказ остался одной строкой. В SQL:

SELECT
    order_id,
    transform(prices, p -> p * 0.9) AS discounted_prices
FROM orders

filter()

Фильтрует элементы массива по условию, возвращает массив только подходящих элементов.

from pyspark.sql.functions import filter

# Оставить только позиции дороже 1000 рублей
data = [(1, [("SKU-A", 500.0), ("SKU-B", 1500.0), ("SKU-C", 200.0)])]
# Для простоты - массив чисел
data = [(1, [500.0, 1500.0, 200.0, 2000.0])]
df = spark.createDataFrame(data, ["order_id", "prices"])

df.select(
    "order_id",
    filter("prices", lambda x: x > 1000).alias("premium_prices")
).show()
+--------+--------------+
|order_id|premium_prices|
+--------+--------------+
|       1| [1500.0, 2000.0]|
+--------+--------------+

aggregate()

Свёртка массива в одно значение (сумма, максимум, строка и т.д.). Аналог reduce / fold.

from pyspark.sql.functions import aggregate

# Сумма всех цен в заказе без explode
df.select(
    "order_id",
    aggregate(
        "prices",
        lit(0.0),       # начальное значение аккумулятора
        lambda acc, x: acc + x  # лямбда: аккумулятор + текущий элемент
    ).alias("total_price")
).show()
+--------+-----------+
|order_id|total_price|
+--------+-----------+
|       1|     4200.0|
+--------+-----------+

В SQL это выглядит так:

SELECT
    order_id,
    aggregate(prices, CAST(0.0 AS DOUBLE), (acc, x) -> acc + x) AS total_price
FROM orders

exists() и forall()

from pyspark.sql.functions import exists, forall

# Есть ли хоть одна цена выше 1000?
df.select("order_id", exists("prices", lambda x: x > 1000).alias("has_premium")).show()

# Все ли цены выше 100?
df.select("order_id", forall("prices", lambda x: x > 100).alias("all_premium")).show()

Когда HOF лучше explode

Используйте higher-order functions вместо explode когда:

Задача Рекомендация
Нормализация (1 строка → N строк) explode
Трансформация каждого элемента transform
Подсчёт суммы/max/min по массиву aggregate
Фильтрация элементов внутри массива filter
Проверка условия на массиве exists / forall
JOIN с другой таблицей по элементам explode + JOIN

Главное правило: если вам не нужно увеличивать количество строк - не делайте этого. Higher-order functions на порядок дешевле: нет GenerateExec, нет дополнительного shuffle из-за изменения размеров партиций.


Производительность и shuffle pressure

Разворачивание массивов влияет на производительность двумя путями.

Рост объёма данных

После explode данных (по числу строк) становится в K раз больше, где K - средний размер массива. Если дальше нужен shuffle (groupBy, join), он будет работать с K-кратно большим датасетом.

# Антипаттерн: explode → groupBy → shuffle с огромным датасетом
df.select("user_id", explode("purchases").alias("purchase")) \
    .groupBy("user_id") \
    .agg(spark_sum("purchase.amount").alias("total_spent"))
# Оптимально: aggregate без explode
from pyspark.sql.functions import aggregate, col

df.select(
    "user_id",
    aggregate(
        "purchases",
        lit(0.0),
        lambda acc, x: acc + x.getField("amount")
    ).alias("total_spent")
)

Второй вариант никогда не умножает строки, shuffle происходит только один раз над исходным датасетом.

Partition skew

Если массивы неравномерны по размеру, после explode возникает skew. Партиции, которые содержали строки с большими массивами, вырастают в разы.

Решение: после explode перераспределить данные если следующий шаг - агрегация:

df_exploded = df.select("user_id", explode("tags").alias("tag"))

# Рекластеризуем по ключу агрегации
df_exploded = df_exploded.repartition(200, "tag")

# Теперь groupBy("tag") не будет дополнительно шафлить
df_exploded.groupBy("tag").count().show()

Почему GenerateExec не вызывает shuffle сам по себе

Это важно понимать: explode и GenerateExec - локальные операции. Они не перемещают данные между узлами. Каждый executor разворачивает свои партиции независимо. Shuffle происходит только если следующая операция его требует (join, groupBy, sort). Но объём данных при этом shuffle уже больше в K раз.


Паттерн Bronze/Silver с unnesting

В production-системах разворачивание массивов - стандартный шаг преобразования из Bronze в Silver.

from pyspark.sql.functions import col, inline_outer, to_date, current_timestamp

# Bronze: читаем из Delta/Parquet как есть
bronze_orders = spark.read.format("delta").load("s3a://datalake/bronze/orders/")

# Silver: нормализованные позиции заказов
silver_order_items = (
    bronze_orders
    .filter(col("_corrupt_record").isNull())  # убираем битые записи
    .select(
        col("order_id"),
        col("customer_id"),
        to_date("created_at").alias("order_date"),
        col("currency"),
        col("status"),
        inline_outer("line_items"),  # разворачиваем позиции
    )
    .withColumn("extracted_at", current_timestamp())
)

# Пишем в Silver с партиционированием по дате
silver_order_items.write \
    .format("delta") \
    .mode("overwrite") \
    .partitionBy("order_date") \
    .option("overwriteSchema", "true") \
    .save("s3a://datalake/silver/order_items/")

Инкрементальная загрузка через Delta MERGE

При инкрементальной обработке нужно учитывать идемпотентность: если заказ обновился (добавилась позиция), нужно обновить Silver-таблицу.

from delta.tables import DeltaTable

silver_table = DeltaTable.forPath(spark, "s3a://datalake/silver/order_items/")

# Новые/изменённые заказы из Bronze
new_batch = (
    bronze_orders
    .filter(col("updated_at") > last_processed_ts)
    .select(col("order_id"), col("customer_id"), inline_outer("line_items"))
)

# MERGE: обновляем существующие, вставляем новые
silver_table.alias("silver").merge(
    new_batch.alias("new"),
    "silver.order_id = new.order_id AND silver.sku = new.sku"
).whenMatchedUpdateAll() \
 .whenNotMatchedInsertAll() \
 .execute()

Работа с несколькими массивами в одной записи

Иногда запись содержит несколько независимых массивов. Правильная стратегия - создать отдельную нормализованную таблицу для каждого:

# Заказы содержат позиции (items) и историю статусов (status_history)
# Создаём две Silver-таблицы

# Silver: позиции
silver_items = bronze_orders.select(
    "order_id", "customer_id",
    inline_outer("items")
).write.format("delta").save("s3a://datalake/silver/order_items/")

# Silver: история статусов (отдельная таблица!)
silver_statuses = bronze_orders.select(
    "order_id",
    posexplode_outer("status_history").alias("step_num", "status_record"),
    col("status_record.status"),
    col("status_record.changed_at"),
    col("status_record.changed_by"),
).write.format("delta").save("s3a://datalake/silver/order_statuses/")

Две отдельные Silver-таблицы - правильная нормализация. Объединять их через Double Explode в одном select - антипаттерн.


Полный пример с LATERAL VIEW в Spark SQL

Покажем полный аналитический запрос с использованием SQL-синтаксиса:

# Регистрируем Bronze-таблицу как temporary view
bronze_orders.createOrReplaceTempView("bronze_orders")

# Аналитика: топ-10 SKU по выручке за последние 30 дней
result = spark.sql("""
    WITH order_items AS (
        SELECT
            o.order_id,
            o.customer_id,
            o.order_date,
            item.sku,
            item.product_name,
            item.category,
            item.qty,
            item.unit_price,
            item.discount_pct,
            ROUND(item.qty * item.unit_price * (1 - item.discount_pct), 2) AS line_total
        FROM bronze_orders o
        LATERAL VIEW OUTER explode(o.items) tmp AS item
        WHERE o.order_date >= DATE_SUB(CURRENT_DATE(), 30)
          AND o.status = 'completed'
    )
    SELECT
        sku,
        product_name,
        category,
        SUM(qty)         AS units_sold,
        SUM(line_total)  AS revenue,
        COUNT(DISTINCT order_id) AS orders_count,
        ROUND(AVG(unit_price), 2) AS avg_unit_price
    FROM order_items
    WHERE sku IS NOT NULL  -- фильтруем null от outer
    GROUP BY sku, product_name, category
    ORDER BY revenue DESC
    LIMIT 10
""")

result.show(truncate=False)

CTE order_items делает разворачивание один раз, дальнейшая аналитика работает с нормализованным результатом. Это читаемо и эффективно: Spark оптимизирует весь запрос через Catalyst, включая predicate pushdown на order_date и status.


Работа с deeply nested структурами

Иногда данные вложены на несколько уровней: массив содержит struct, который содержит массив.

# Пример: пользователи → заказы → позиции (три уровня)
deep_data = [(
    1,  # user_id
    [   # orders[]
        {
            "order_id": 101,
            "items": [
                {"sku": "A", "qty": 2},
                {"sku": "B", "qty": 1},
            ]
        },
        {
            "order_id": 102,
            "items": [
                {"sku": "C", "qty": 3},
            ]
        }
    ]
)]

Разворачивание происходит поэтапно - по одному уровню за раз:

from pyspark.sql.functions import explode, inline

# Шаг 1: разворачиваем заказы
df_orders = df.select(
    col("user_id"),
    explode("orders").alias("order")
).select(
    col("user_id"),
    col("order.order_id"),
    col("order.items"),
)

# Шаг 2: разворачиваем позиции
df_items = df_orders.select(
    col("user_id"),
    col("order_id"),
    inline("items"),  # разворачивает struct-поля
)

df_items.show()
+-------+--------+---+---+
|user_id|order_id|sku|qty|
+-------+--------+---+---+
|      1|     101|  A|  2|
|      1|     101|  B|  1|
|      1|     102|  C|  3|
+-------+--------+---+---+

Практические задачи

Задача 1: Анализ тегов статей

data = [
    (1, "Intro to Spark",    ["spark", "python", "data-engineering"]),
    (2, "Kafka Streams",     ["kafka", "streaming", "java"]),
    (3, "Delta Lake Guide",  ["delta", "spark", "lakehouse"]),
    (4, "Python Basics",     ["python", "beginner"]),
    (5, "Empty Article",     []),
]

articles = spark.createDataFrame(data, ["article_id", "title", "tags"])

# 1. Топ-5 тегов по частоте использования
from pyspark.sql.functions import explode_outer, count, desc

articles \
    .select("article_id", explode_outer("tags").alias("tag")) \
    .filter(col("tag").isNotNull()) \
    .groupBy("tag") \
    .agg(count("article_id").alias("article_count")) \
    .orderBy(desc("article_count")) \
    .limit(5) \
    .show()

# 2. Статьи с тегом "spark"
articles \
    .select("article_id", "title", explode("tags").alias("tag")) \
    .filter(col("tag") == "spark") \
    .select("article_id", "title") \
    .show()

# 3. Статьи без тегов (используем explode_outer)
articles \
    .select("article_id", "title", explode_outer("tags").alias("tag")) \
    .filter(col("tag").isNull()) \
    .select("article_id", "title") \
    .show()

Задача 2: Воронка конверсии по шагам

data = [
    (1, ["landing", "product", "cart", "checkout", "purchase"]),
    (2, ["landing", "product", "cart"]),
    (3, ["landing", "product"]),
    (4, ["landing", "product", "cart", "checkout"]),
]

sessions = spark.createDataFrame(data, ["session_id", "funnel"])

from pyspark.sql.functions import posexplode, count, countDistinct

# Сколько сессий доходит до каждого шага?
sessions \
    .select("session_id", posexplode("funnel").alias("step_num", "step")) \
    .groupBy("step_num", "step") \
    .agg(countDistinct("session_id").alias("sessions_reached")) \
    .orderBy("step_num") \
    .show()
+--------+----------+----------------+
|step_num|      step|sessions_reached|
+--------+----------+----------------+
|       0|   landing|               4|
|       1|   product|               4|
|       2|      cart|               3|
|       3|  checkout|               2|
|       4|  purchase|               1|
+--------+----------+----------------+

Конверсия от landing до purchase: 1/4 = 25%. Самый большой отток - между cart и checkout.

Задача 3: Парные теги (co-occurrence)

Найти пары тегов, которые часто встречаются вместе:

# Self-join после explode для нахождения пар
tags_df = articles.select("article_id", explode("tags").alias("tag"))

# Cross-join с самим собой, исключаем пары tag=tag и дубликаты
tag_pairs = tags_df.alias("t1").join(
    tags_df.alias("t2"),
    (col("t1.article_id") == col("t2.article_id")) &
    (col("t1.tag") < col("t2.tag"))  # < чтобы избежать (A,B) и (B,A)
).select(
    col("t1.tag").alias("tag1"),
    col("t2.tag").alias("tag2")
)

tag_pairs.groupBy("tag1", "tag2") \
    .agg(count("*").alias("co_occurrences")) \
    .orderBy(desc("co_occurrences")) \
    .show()

Чеклист

  • [ ] Знаю разницу между explode и explode_outer (поведение с null/пустыми массивами)
  • [ ] Использую inline / inline_outer для array<struct<...>> вместо explode + .field
  • [ ] Понимаю зачем нужен posexplode (порядок элементов важен)
  • [ ] Знаю что такое GenerateExec в физическом плане и почему он не вызывает shuffle сам по себе
  • [ ] Не применяю Double Explode к независимым массивам (понимаю что это декартово произведение)
  • [ ] Оцениваю cardinality explosion перед explode больших массивов (percentile_approx на size)
  • [ ] Знаю three higher-order functions: transform, filter, aggregate - и когда они лучше explode
  • [ ] Применяю unnesting только в Silver-слое, Bronze хранит нативную вложенность
  • [ ] Могу написать LATERAL VIEW синтаксис в Spark SQL с несколькими VIEW
  • [ ] Понимаю partition skew после explode и знаю как его решить (repartition после)