Complex Types: работа с ArrayType, MapType и функциями higher-order

ArrayType, MapType, StructType, explode, higher-order functions (transform, filter, aggregate) - обработка вложенных данных без взрыва строк

core optimization

Зачем вообще нужны сложные типы

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

Представьте типичный сценарий: пользователь делает заказ в интернет-магазине и покупает несколько товаров одновременно. В реляционной модели для этого нужны минимум три таблицы: orders, order_items, products. В lakehouse-архитектуре Bronze-слой хранит события такими, какими они приходят - из Kafka, REST API, мобильного приложения - то есть в виде JSON-документа, где список товаров лежит прямо внутри события заказа.

{
  "order_id": "ORD-1001",
  "user_id": 42,
  "items": [
    {"product_id": "P-100", "name": "Ноутбук", "price": 89999, "qty": 1},
    {"product_id": "P-205", "name": "Мышь", "price": 1499, "qty": 2}
  ],
  "metadata": {"platform": "ios", "app_version": "3.4.1", "region": "RU"}
}

Spark поддерживает хранение и обработку таких структур нативно, без предварительного разворачивания в отдельные строки. Это и есть Complex Types - сложные типы данных.

Важно понимать: Complex Types - это не «костыль» и не временное решение. В правильно спроектированной lakehouse-архитектуре:

  • Bronze-слой хранит сырые данные в том виде, в котором они пришли - с вложенностью, массивами, динамическими ключами
  • Silver-слой делает выборочное уплощение (flattening) только тех полей, которые нужны для аналитики, оставляя вложенность там, где она уместна
  • Gold-слой (витрины) уже работает с плоскими, денормализованными агрегатами

Умение работать со сложными типами - ключевой навык для Data Engineer, потому что он позволяет избежать преждевременной нормализации и работать с данными эффективно на каждом слое.

Три сложных типа в Spark

Spark поддерживает три вида сложных типов. Каждый решает свою задачу:

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

MapType - словарь ключ → значение с произвольным набором ключей. Подходит, когда набор ключей заранее неизвестен или меняется: метаданные события, пользовательские атрибуты, динамические метки. В Parquet и ORC хранится эффективно как пара колонок (keys array + values array).

StructType - именованная структура с фиксированным набором полей и типов, аналог dataclass в Python или record в других языках. Используется для вложенных объектов с известной схемой: адрес, ценовая структура, геокоординаты. Подробно рассматривается в следующем уроке; здесь мы коснёмся только базового доступа к полям.

ArrayType: массивы внутри DataFrame

Как выглядит ArrayType в схеме

Когда Spark читает JSON с массивами, он автоматически определяет тип как ArrayType(elementType):

from pyspark.sql import SparkSession
from pyspark.sql.types import (
    StructType, StructField, StringType, IntegerType,
    DoubleType, ArrayType, MapType
)

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

# Данные заказов: каждый заказ содержит массив товаров
data = [
    ("ORD-1001", 42, ["laptop", "mouse", "keyboard"], [89999.0, 1499.0, 3499.0]),
    ("ORD-1002", 17, ["headphones"],                  [7990.0]),
    ("ORD-1003", 99, ["tablet", "case"],               [45000.0, 890.0]),
]

schema = StructType([
    StructField("order_id",  StringType(),                   False),
    StructField("user_id",   IntegerType(),                  False),
    StructField("products",  ArrayType(StringType()),         True),
    StructField("prices",    ArrayType(DoubleType()),         True),
])

df = spark.createDataFrame(data, schema)
df.printSchema()
# root
#  |-- order_id: string (nullable = false)
#  |-- user_id: integer (nullable = false)
#  |-- products: array (nullable = true)
#  |    |-- element: string (containsNull = true)
#  |-- prices: array (nullable = true)
#  |    |-- element: double (containsNull = true)

Обратите внимание на containsNull = true у элементов массива - это отдельный флаг, означающий, что сам элемент в массиве может быть null, независимо от того, nullable ли вся колонка-массив.

Создание массивов в SQL и API

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

from pyspark.sql.functions import array, col, lit

# Способ 1: функция array() - объединяет несколько колонок в массив
flat_df = spark.createDataFrame([
    ("P-100", "laptop",  89999.0, 0.1),
    ("P-205", "mouse",    1499.0, 0.0),
], ["product_id", "name", "price", "discount"])

# Создаём массив из двух числовых колонок
flat_df.select(
    col("product_id"),
    array(col("price"), col("price") * (1 - col("discount"))).alias("price_range")
).show()
# +----------+-------------------+
# |product_id|        price_range|
# +----------+-------------------+
# |     P-100|[89999.0, 80999.1] |
# |     P-205|  [1499.0, 1499.0] |
# +----------+-------------------+

# Способ 2: collect_list() при агрегации - собирает значения из группы в массив
from pyspark.sql.functions import collect_list, collect_set

items_df = spark.createDataFrame([
    ("ORD-1001", "laptop"),
    ("ORD-1001", "mouse"),
    ("ORD-1001", "mouse"),  # дубликат
    ("ORD-1002", "tablet"),
], ["order_id", "product"])

items_df.groupBy("order_id").agg(
    collect_list("product").alias("products_with_duplicates"),  # сохраняет дубликаты
    collect_set("product").alias("unique_products"),             # убирает дубликаты
).show(truncate=False)
# +---------+------------------------+---------------+
# |order_id |products_with_duplicates|unique_products|
# +---------+------------------------+---------------+
# |ORD-1001 |[laptop, mouse, mouse]  |[laptop, mouse]|
# |ORD-1002 |[tablet]                |[tablet]       |
# +---------+------------------------+---------------+

collect_list vs collect_set - ключевое различие: collect_list сохраняет порядок и дубликаты (результат зависит от порядка обработки partition, поэтому порядок может быть непредсказуемым), collect_set возвращает уникальные элементы без гарантии порядка.

Доступ к элементам массива

Spark поддерживает несколько способов обратиться к конкретному элементу:

from pyspark.sql.functions import element_at

# Индексация через [] - нумерация с 0 (стандарт Python)
df.select(
    col("order_id"),
    col("products")[0].alias("first_product"),    # первый элемент
    col("products")[1].alias("second_product"),   # второй элемент (None если нет)
).show()

# element_at() - нумерация с 1, поддерживает отрицательные индексы
df.select(
    col("order_id"),
    element_at(col("products"), 1).alias("first_product"),   # первый (индекс 1)
    element_at(col("products"), -1).alias("last_product"),   # последний (индекс -1)
).show()

Почему два разных подхода? col[idx] унаследован от Scala-стиля с 0-индексацией. element_at() добавлен позже и следует SQL-конвенции (индексация с 1), плюс поддерживает отрицательные индексы для обращения с конца массива. В новом коде рекомендуется element_at() - он явнее и безопаснее.

Базовые функции для работы с массивами

from pyspark.sql.functions import (
    size, array_contains, sort_array, array_distinct,
    array_union, array_intersect, array_except, flatten
)

# size() - количество элементов в массиве (null → -1 по умолчанию)
df.select(
    col("order_id"),
    size(col("products")).alias("item_count")
).show()

# array_contains() - проверка наличия элемента
df.filter(
    array_contains(col("products"), "laptop")
).select("order_id").show()

# sort_array() - сортировка (по умолчанию ascending=True)
df.select(
    col("order_id"),
    sort_array(col("prices"), asc=False).alias("prices_desc")
).show()

# array_distinct() - убирает дубликаты, сохраняет первое вхождение
# array_union() - объединение двух массивов (без дубликатов)
# array_intersect() - пересечение (общие элементы)
# array_except() - разность (элементы первого, которых нет во втором)

# Пример: какие продукты есть в заказе, но нет в wishlist пользователя?
orders_and_wishes = spark.createDataFrame([
    ("ORD-1001", ["laptop", "mouse", "keyboard"], ["keyboard", "webcam"]),
], ["order_id", "ordered", "wishlist"])

orders_and_wishes.select(
    col("order_id"),
    array_intersect(col("ordered"), col("wishlist")).alias("from_wishlist"),
    array_except(col("ordered"), col("wishlist")).alias("impulsive_purchases"),
).show(truncate=False)
# +---------+-------------+--------------------+
# |order_id |from_wishlist|impulsive_purchases |
# +---------+-------------+--------------------+
# |ORD-1001 |[keyboard]   |[laptop, mouse]     |
# +---------+-------------+--------------------+

# flatten() - «схлопывает» массив массивов в один плоский массив
nested = spark.createDataFrame([
    (1, [["a", "b"], ["c"]]),
    (2, [["x"], ["y", "z"]]),
], ["id", "nested_array"])

nested.select(
    col("id"),
    flatten(col("nested_array")).alias("flat_array")
).show()
# +---+----------+
# | id|flat_array|
# +---+----------+
# |  1|[a, b, c] |
# |  2|[x, y, z] |
# +---+----------+

explode: разворачивание массива в строки

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

explode() преобразует каждый элемент массива в отдельную строку. Это мощный инструмент, но его нужно использовать осознанно.

Заметьте: колонки order_id и user_id дублируются в каждой строке. Если массив содержит 100 элементов, строка размножится в 100 строк, и все остальные колонки будут скопированы 100 раз. Это называется data explosion - именно отсюда происходит название функции.

from pyspark.sql.functions import explode, posexplode, explode_outer

# explode() - каждый элемент становится строкой
# Если массив пустой или null - строка пропадает из результата
df.select(
    col("order_id"),
    col("user_id"),
    explode(col("products")).alias("product")
).show()
# +---------+-------+--------+
# |order_id |user_id|product |
# +---------+-------+--------+
# |ORD-1001 |42     |laptop  |
# |ORD-1001 |42     |mouse   |
# |ORD-1001 |42     |keyboard|
# |ORD-1002 |17     |headphones|
# |ORD-1003 |99     |tablet  |
# |ORD-1003 |99     |case    |
# +---------+-------+--------+

# posexplode() - добавляет колонку с позицией (индексом) элемента
df.select(
    col("order_id"),
    posexplode(col("products")).alias("pos", "product")
).show()
# +---------+---+--------+
# |order_id |pos|product |
# +---------+---+--------+
# |ORD-1001 |0  |laptop  |
# |ORD-1001 |1  |mouse   |
# |ORD-1001 |2  |keyboard|
# ...

# explode_outer() - в отличие от explode, сохраняет строки с пустыми/null массивами
# Это аналог LEFT JOIN по поведению
data_with_empty = spark.createDataFrame([
    ("ORD-1004", 55, []),           # пустой массив
    ("ORD-1005", 66, None),         # null
    ("ORD-1006", 77, ["tablet"]),   # обычный
], ["order_id", "user_id", "products"])

data_with_empty.select(
    col("order_id"),
    explode_outer(col("products")).alias("product")  # null строки сохранятся
).show()
# +---------+-------+
# |order_id |product|
# +---------+-------+
# |ORD-1004 |null   |   ← пустой массив → null
# |ORD-1005 |null   |   ← null массив → null
# |ORD-1006 |tablet |
# +---------+-------+

Почему explode - это дорого

explode выглядит простым, но у него есть серьёзные последствия для производительности:

Умножение строк увеличивает объём shuffle. Если после explode идёт groupBy или join, Spark шафлит уже умноженные строки. Например: 10 миллионов заказов с 20 товарами каждый → после explode → 200 миллионов строк → groupBy → 200 миллионов строк по сети. Тогда как без explode можно агрегировать прямо внутри массива.

Память partition растёт. Partition, которая до explode весила 128 МБ, может стать 2 ГБ после explode массива с 20 элементами. Это давление на Execution Memory и риск spill.

Дублирование данных из других колонок. Каждый раз, когда вы делаете explode, все остальные колонки строки копируются. Если строка содержит тяжёлые колонки (длинный текст, JSON), умножение происходит и для них.

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

Higher-Order Functions: обработка массивов без разворачивания

Higher-Order Functions (функции высшего порядка) - это функции, которые принимают другую функцию (лямбду) как аргумент. В контексте Spark это означает: обработка каждого элемента массива прямо внутри строки, без explode, без умножения строк, без дополнительного shuffle.

Как Catalyst исполняет Higher-Order Functions

Важно понять: Higher-Order Functions компилируются Catalyst в Java-код через WholeStageCodegen. Это означает, что лямбда-выражение, которое вы пишете на Python, никогда не выполняется в Python - оно транслируется в JVM-байткод и выполняется напрямую в JVM. Никакого IPC с Python worker'ом, никакой сериализации через pickle. По производительности это эквивалентно встроенным Spark-функциям.

transform(): применить функцию к каждому элементу

transform() - аналог map() в Python. Принимает массив и лямбду, возвращает новый массив той же длины, где каждый элемент преобразован.

from pyspark.sql.functions import transform, col, upper, round as spark_round

# Лямбда обозначается как x -> expression
# В Python API: lambda x: expression

# Пример 1: привести все названия продуктов к верхнему регистру
df.select(
    col("order_id"),
    transform(col("products"), lambda x: upper(x)).alias("products_upper")
).show(truncate=False)
# +---------+----------------------------+
# |order_id |products_upper              |
# +---------+----------------------------+
# |ORD-1001 |[LAPTOP, MOUSE, KEYBOARD]   |
# |ORD-1002 |[HEADPHONES]                |
# |ORD-1003 |[TABLET, CASE]              |
# +---------+----------------------------+

# Пример 2: применить скидку 10% ко всем ценам
df.select(
    col("order_id"),
    transform(
        col("prices"),
        lambda p: spark_round(p * 0.9, 2)
    ).alias("discounted_prices")
).show(truncate=False)
# +---------+----------------------------+
# |order_id |discounted_prices           |
# +---------+----------------------------+
# |ORD-1001 |[80999.1, 1349.1, 3149.1]  |

# Пример 3: transform с SQL-синтаксисом через spark.sql
# SQL-синтаксис: TRANSFORM(array, element -> expression)
spark.sql("""
    SELECT order_id,
           TRANSFORM(prices, p -> ROUND(p * 0.9, 2)) AS discounted_prices
    FROM orders_view
""")

transform() не меняет длину массива - входной массив из N элементов всегда даёт выходной массив из N элементов. Если лямбда вернула null для какого-то элемента - в этой позиции будет null.

filter(): оставить только подходящие элементы

filter() применяет предикат к каждому элементу и возвращает массив только из тех элементов, для которых предикат вернул true. Длина результирующего массива может быть меньше исходного.

from pyspark.sql.functions import filter as spark_filter

# Данные с ценами и флагом доступности товаров
orders_stock = spark.createDataFrame([
    ("ORD-1001",
     ["laptop", "mouse", "keyboard"],
     [89999.0, 1499.0, 3499.0],
     [True, True, False]),   # keyboard нет в наличии
    ("ORD-1002",
     ["headphones", "stand"],
     [7990.0, 2500.0],
     [False, True]),          # headphones нет
], ["order_id", "products", "prices", "in_stock"])

# filter() с одним аргументом - только элемент
# Предположим: хотим оставить только дорогие товары (цена > 5000)
orders_stock.select(
    col("order_id"),
    spark_filter(
        col("prices"),
        lambda p: p > 5000
    ).alias("expensive_prices")
).show(truncate=False)
# +---------+------------------+
# |order_id |expensive_prices  |
# +---------+------------------+
# |ORD-1001 |[89999.0]         |
# |ORD-1002 |[7990.0]          |
# +---------+------------------+

# filter() с двумя аргументами - элемент и индекс
# Полезно когда нужно учесть позицию: оставить только чётные позиции
orders_stock.select(
    col("order_id"),
    spark_filter(
        col("products"),
        lambda x, i: i % 2 == 0   # i - индекс (0, 1, 2...)
    ).alias("even_indexed_products")
).show(truncate=False)
# +---------+------------------------+
# |order_id |even_indexed_products   |
# +---------+------------------------+
# |ORD-1001 |[laptop, keyboard]      |   ← позиции 0, 2
# |ORD-1002 |[headphones]            |   ← позиция 0
# +---------+------------------------+

Обратите внимание: в Python API filter - встроенная функция, поэтому импортируем как spark_filter. В SQL синтаксе конфликта нет: FILTER(array, element -> condition).

exists() и forall(): проверка условий на коллекции

exists() проверяет, удовлетворяет ли хотя бы один элемент условию (аналог any() в Python). forall() проверяет, удовлетворяют ли все элементы условию (аналог all()). Обе возвращают Boolean-колонку, а не массив.

from pyspark.sql.functions import exists, forall

# exists(): есть ли в заказе хотя бы один товар дороже 50 000?
orders_stock.select(
    col("order_id"),
    exists(
        col("prices"),
        lambda p: p > 50000
    ).alias("has_premium_item")
).show()
# +---------+----------------+
# |order_id |has_premium_item|
# +---------+----------------+
# |ORD-1001 |true            |   ← laptop = 89999
# |ORD-1002 |false           |   ← нет товаров > 50000
# +---------+----------------+

# forall(): все ли товары в заказе стоят меньше 10 000?
orders_stock.select(
    col("order_id"),
    forall(
        col("prices"),
        lambda p: p < 10000
    ).alias("all_budget_items")
).show()
# +---------+----------------+
# |order_id |all_budget_items|
# +---------+----------------+
# |ORD-1001 |false           |   ← laptop = 89999 нарушает условие
# |ORD-1002 |true            |   ← 7990 и 2500, оба < 10000
# +---------+----------------+

# Комбинирование: использовать exists/forall как фильтр DataFrame
# Оставить только заказы с хотя бы одним товаром > 5000
orders_stock.filter(
    exists(col("prices"), lambda p: p > 5000)
).select("order_id").show()

exists() и forall() - семантически эквивалентны SQL EXISTS (subquery) и ALL (subquery), но работают над массивом внутри строки. Это особенно удобно для фильтрации DataFrame по условиям на вложенных коллекциях без разворачивания.

aggregate(): свернуть массив в одно значение

aggregate() (он же reduce() в функциональном программировании) - самая мощная из Higher-Order Functions. Она проходит по массиву, накапливая промежуточное значение (аккумулятор), и возвращает финальный результат.

Сигнатура: aggregate(array, initial_value, merge_func, finish_func=None)

  • initial_value - начальное значение аккумулятора
  • merge_func(accumulator, element) -> accumulator - как обновить аккумулятор следующим элементом
  • finish_func(accumulator) -> result - финальное преобразование (необязательно)
from pyspark.sql.functions import aggregate, lit

# Пример 1: простая сумма цен (то же что array_sum, но демонстрирует механику)
df.select(
    col("order_id"),
    aggregate(
        col("prices"),
        lit(0.0),                              # начальный аккумулятор
        lambda acc, p: acc + p                 # acc + каждый элемент
    ).alias("total_price")
).show()
# +---------+-----------+
# |order_id |total_price|
# +---------+-----------+
# |ORD-1001 |94997.0    |   ← 89999 + 1499 + 3499
# |ORD-1002 |7990.0     |
# |ORD-1003 |45890.0    |   ← 45000 + 890

# Пример 2: сложная логика - скидка 15% только на товары дороже 5000
# Это то, что нельзя сделать одной встроенной функцией
df.select(
    col("order_id"),
    aggregate(
        col("prices"),
        lit(0.0),
        lambda acc, p: acc + (
            spark_round(p * 0.85, 2)  # скидка 15% для дорогих
            if p > 5000
            else p                    # без скидки для дешёвых
        )
        # Но в Spark лямбде нет if/else Python - нужен when/otherwise:
    ).alias("total_with_selective_discount")
)

# Правильная версия с when/otherwise внутри лямбды:
from pyspark.sql.functions import when

df.select(
    col("order_id"),
    aggregate(
        col("prices"),
        lit(0.0),
        lambda acc, p: acc + when(p > 5000, spark_round(p * 0.85, 2)).otherwise(p)
    ).alias("total_selective_discount")
).show()
# +---------+------------------------+
# |order_id |total_selective_discount|
# +---------+------------------------+
# |ORD-1001 |79149.7                 |   ← 89999×0.85 + 1499 + 3499×0.85
# |ORD-1002 |6791.5                  |   ← 7990×0.85
# |ORD-1003 |39140.0                 |   ← 45000×0.85 + 890

Ключевой момент: внутри лямбды Higher-Order Functions нельзя использовать Python-операторы if/else - только Spark-выражения (when/otherwise, арифметику, функции F.*). Это потому что лямбда транслируется в JVM-код, а Python if остаётся в Python и не видим Catalyst.

zip_with(): параллельная обработка двух массивов

zip_with() объединяет два массива поэлементно, применяя бинарную функцию. Первый элемент первого массива + первый элемент второго → первый результат. Массивы должны быть одинаковой длины.

from pyspark.sql.functions import zip_with

# Вычислить стоимость позиций: цена × количество
orders_with_qty = spark.createDataFrame([
    ("ORD-1001", ["laptop", "mouse"], [89999.0, 1499.0], [1, 2]),
    ("ORD-1002", ["tablet", "case"],  [45000.0, 890.0],  [1, 3]),
], ["order_id", "products", "prices", "quantities"])

orders_with_qty.select(
    col("order_id"),
    zip_with(
        col("prices"),
        col("quantities"),
        lambda price, qty: spark_round(price * qty, 2)
    ).alias("line_totals")
).show(truncate=False)
# +---------+------------------+
# |order_id |line_totals       |
# +---------+------------------+
# |ORD-1001 |[89999.0, 2998.0] |   ← 89999×1, 1499×2
# |ORD-1002 |[45000.0, 2670.0] |   ← 45000×1, 890×3
# +---------+------------------+

# Комбинация zip_with + aggregate: итоговая сумма заказа
orders_with_qty.select(
    col("order_id"),
    aggregate(
        zip_with(col("prices"), col("quantities"), lambda p, q: p * q),
        lit(0.0),
        lambda acc, x: acc + x
    ).alias("order_total")
).show()
# +---------+-----------+
# |order_id |order_total|
# +---------+-----------+
# |ORD-1001 |92997.0    |   ← 89999 + 2998
# |ORD-1002 |47670.0    |   ← 45000 + 2670
# +---------+-----------+

zip_with + aggregate - мощная комбинация для вычисления агрегатов по параллельным массивам без единого explode.

MapType: словари внутри DataFrame

Что такое MapType и когда его использовать

MapType(keyType, valueType) - это структура ключ→значение, где набор ключей не фиксирован в схеме. В отличие от StructType, где каждое поле известно заранее, Map хранит произвольный набор пар. Это идеально для:

  • Метаданных события: {"platform": "ios", "version": "3.4.1", "region": "RU"}
  • Пользовательских атрибутов: разные пользователи могут иметь разные наборы атрибутов
  • Конфигурационных параметров: словари настроек с непредсказуемым набором ключей
  • Счётчиков событий: {"click": 5, "view": 12, "purchase": 1}

Создание MapType

from pyspark.sql.functions import create_map, map_from_arrays, map_from_entries

# Способ 1: create_map() - из перечисленных ключей и значений
flat_df.select(
    col("product_id"),
    create_map(
        lit("price"),   col("price"),       # ключ, значение
        lit("discount"), col("discount"),   # ключ, значение
    ).alias("attributes")
).show(truncate=False)
# +----------+---------------------------+
# |product_id|attributes                 |
# +----------+---------------------------+
# |P-100     |{price -> 89999.0, ...}    |

# Способ 2: map_from_arrays() - из двух массивов (ключи и значения)
keys_values_df = spark.createDataFrame([
    ("EVT-001", ["platform", "version", "region"], ["ios", "3.4.1", "RU"]),
    ("EVT-002", ["platform", "region"],             ["android", "DE"]),
], ["event_id", "keys", "values"])

keys_values_df.select(
    col("event_id"),
    map_from_arrays(col("keys"), col("values")).alias("metadata")
).show(truncate=False)
# +--------+-------------------------------------------+
# |event_id|metadata                                   |
# +--------+-------------------------------------------+
# |EVT-001 |{platform -> ios, version -> 3.4.1, region -> RU}|
# |EVT-002 |{platform -> android, region -> DE}        |
# +--------+-------------------------------------------+

# Способ 3: прямые данные в createDataFrame со схемой
events_data = [
    ("EVT-001", {"platform": "ios", "version": "3.4.1"}),
    ("EVT-002", {"platform": "android", "ab_group": "B"}),
]
events_schema = StructType([
    StructField("event_id",  StringType()),
    StructField("metadata",  MapType(StringType(), StringType())),
])
events_df = spark.createDataFrame(events_data, events_schema)

Доступ к значениям Map

from pyspark.sql.functions import map_keys, map_values, element_at, map_contains_key

# Получить значение по ключу через []
events_df.select(
    col("event_id"),
    col("metadata")["platform"].alias("platform"),    # None если ключа нет
    col("metadata")["version"].alias("app_version"),  # None если ключа нет
).show()

# element_at() - то же, но явнее
events_df.select(
    element_at(col("metadata"), "platform").alias("platform")
)

# map_keys() - все ключи в виде массива
events_df.select(
    col("event_id"),
    map_keys(col("metadata")).alias("available_keys")
).show(truncate=False)
# +--------+-----------------------------+
# |event_id|available_keys               |
# +--------+-----------------------------+
# |EVT-001 |[platform, version]          |
# |EVT-002 |[platform, ab_group]         |
# +--------+-----------------------------+

# map_values() - все значения в виде массива
events_df.select(
    col("event_id"),
    map_values(col("metadata")).alias("all_values")
)

# map_contains_key() - проверить наличие ключа (Spark 3.3+)
events_df.filter(
    map_contains_key(col("metadata"), "ab_group")
).select("event_id").show()
# +--------+
# |event_id|
# +--------+
# |EVT-002 |   ← только у него есть ab_group
# +--------+

Операции над Map

from pyspark.sql.functions import map_concat, map_filter, transform_values, transform_keys

# map_concat() - объединить два словаря (ключи из второго перезаписывают первый)
base_meta = spark.createDataFrame([
    ("EVT-001", {"env": "prod", "region": "RU"}, {"version": "3.4.1", "region": "US"}),
], ["event_id", "base", "override"])

base_meta.select(
    col("event_id"),
    map_concat(col("base"), col("override")).alias("merged")  # US перезапишет RU
).show(truncate=False)
# +--------+-------------------------------------------+
# |event_id|merged                                     |
# +--------+-------------------------------------------+
# |EVT-001 |{env -> prod, region -> US, version -> 3.4.1}|
# +--------+-------------------------------------------+

# transform_values() - применить функцию ко всем значениям (ключи не меняются)
# Пример: добавить суффикс к каждому значению metadata
events_df.select(
    col("event_id"),
    transform_values(
        col("metadata"),
        lambda k, v: upper(v)   # v - значение, k - ключ (доступен, но необязателен)
    ).alias("metadata_upper")
).show(truncate=False)
# Все значения приведены к верхнему регистру

# transform_keys() - применить функцию к ключам (значения не меняются)
events_df.select(
    transform_keys(
        col("metadata"),
        lambda k, v: upper(k)   # k - ключ, v - значение
    ).alias("upper_keys_meta")
)

# map_filter() - оставить только пары, удовлетворяющие предикату
events_df.select(
    col("event_id"),
    map_filter(
        col("metadata"),
        lambda k, v: k != "version"   # убрать пары с ключом "version"
    ).alias("metadata_no_version")
)

Разворачивание Map через explode

Как и ArrayType, MapType можно развернуть в строки - explode на Map даёт две колонки: key и value:

# explode Map → пара (key, value) для каждой строки
events_df.select(
    col("event_id"),
    explode(col("metadata")).alias("key", "value")
).show()
# +--------+--------+--------+
# |event_id|key     |value   |
# +--------+--------+--------+
# |EVT-001 |platform|ios     |
# |EVT-001 |version |3.4.1   |
# |EVT-002 |platform|android |
# |EVT-002 |ab_group|B       |
# +--------+--------+--------+

Это полезно когда нужно сделать pivot или агрегацию по динамическим ключам. Но как и с ArrayType, это умножение строк - применяйте осознанно.

StructType: именованные вложенные структуры

StructType - самый распространённый сложный тип в схемах. Когда Spark читает JSON-объект, вложенный объект автоматически становится StructType. В отличие от MapType, набор полей фиксирован в схеме и известен Catalyst - это даёт возможность Nested Column Pruning.

# Данные с вложенной структурой адреса
from pyspark.sql.types import StructField

address_schema = StructType([
    StructField("order_id", StringType()),
    StructField("shipping_address", StructType([
        StructField("city",    StringType()),
        StructField("country", StringType()),
        StructField("zip",     StringType()),
    ])),
    StructField("billing_address", StructType([
        StructField("city",    StringType()),
        StructField("country", StringType()),
    ])),
])

address_data = [
    ("ORD-1001", ("Moscow",     "RU", "101000"), ("Moscow",     "RU")),
    ("ORD-1002", ("Berlin",     "DE", "10115"),  ("Frankfurt",  "DE")),
    ("ORD-1003", ("Amsterdam",  "NL", "1012 AB"),("Amsterdam",  "NL")),
]

addr_df = spark.createDataFrame(address_data, address_schema)
addr_df.printSchema()
# root
#  |-- order_id: string
#  |-- shipping_address: struct
#  |    |-- city: string
#  |    |-- country: string
#  |    |-- zip: string
#  |-- billing_address: struct
#  |    |-- city: string
#  |    |-- country: string

Dot notation: доступ к вложенным полям

Самый простой способ достать поле из структуры - через точку в строковом имени колонки:

# Обращение к вложенному полю через точку в строке
addr_df.select(
    col("order_id"),
    col("shipping_address.city").alias("ship_city"),
    col("shipping_address.country").alias("ship_country"),
    col("billing_address.city").alias("bill_city"),
).show()
# +---------+---------+------------+---------+
# |order_id |ship_city|ship_country|bill_city|
# +---------+---------+------------+---------+
# |ORD-1001 |Moscow   |RU          |Moscow   |
# |ORD-1002 |Berlin   |DE          |Frankfurt|
# |ORD-1003 |Amsterdam|NL          |Amsterdam|
# +---------+---------+------------+---------+

# SQL-синтаксис: то же через SELECT
addr_df.createOrReplaceTempView("orders_addr")
spark.sql("""
    SELECT order_id,
           shipping_address.city AS ship_city,
           billing_address.city  AS bill_city
    FROM orders_addr
    WHERE shipping_address.country = 'RU'
""").show()

# Доступ через getField() - программный вариант когда имя поля в переменной
field_name = "city"
addr_df.select(
    col("shipping_address").getField(field_name).alias(field_name)
)

struct(): создание именованной структуры

from pyspark.sql.functions import struct

# Упаковать плоские колонки в структуру
flat_addresses = spark.createDataFrame([
    ("ORD-1001", "Moscow", "RU", "101000"),
    ("ORD-1002", "Berlin", "DE", "10115"),
], ["order_id", "city", "country", "zip"])

flat_addresses.select(
    col("order_id"),
    struct(
        col("city"),
        col("country"),
        col("zip")
    ).alias("address")
).printSchema()
# root
#  |-- order_id: string
#  |-- address: struct
#  |    |-- city: string
#  |    |-- country: string
#  |    |-- zip: string

Подробно работа с StructType, withField(), schema evolution и глубокая вложенность рассматриваются в следующем уроке.

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

Nested Column Pruning: читаем только нужные поля

Одно из главных преимуществ StructType перед MapType - Catalyst умеет делать Nested Column Pruning: читать из Parquet только те вложенные поля, которые запрашивает запрос, пропуская остальные.

Это работает только для StructType. MapType хранится как две колонки (keys array + values array), и Parquet не умеет выбрать только нужные ключи - читается весь Map. Ещё одна причина предпочитать StructType когда схема известна заранее.

# Проверить в explain: Catalyst должен показать только нужные поля
addr_df.select("order_id", "shipping_address.city").explain("formatted")
# == Physical Plan ==
# *(1) Project [order_id#0, shipping_address#1.city AS city#10]
# +- *(1) Scan parquet [order_id#0, shipping_address.city#2]
#                       ↑                    ↑
#               читает только эти две колонки, не весь struct

Higher-Order Functions vs Python UDF: измеримая разница

import time

# Большой датасет для демонстрации разницы
large_df = spark.range(2_000_000) \
    .withColumn("prices", array(
        (col("id") % 1000 + 1).cast("double"),
        (col("id") % 500 + 1).cast("double"),
        (col("id") % 200 + 1).cast("double"),
    ))

# Способ 1: Python UDF (медленно - IPC на каждую строку)
from pyspark.sql.functions import udf
from pyspark.sql.types import DoubleType

@udf(DoubleType())
def sum_with_discount_udf(prices):
    if prices is None:
        return None
    return sum(p * 0.9 if p > 500 else p for p in prices)

t0 = time.time()
large_df.withColumn("total", sum_with_discount_udf(col("prices"))).count()
print(f"Python UDF: {time.time() - t0:.1f} сек")
# Python UDF: ~45 сек

# Способ 2: Higher-Order Function (быстро - JVM, нет IPC)
t0 = time.time()
large_df.withColumn("total",
    aggregate(
        col("prices"),
        lit(0.0),
        lambda acc, p: acc + when(p > 500, p * 0.9).otherwise(p)
    )
).count()
print(f"Higher-Order Function: {time.time() - t0:.1f} сек")
# Higher-Order Function: ~4 сек - в 10+ раз быстрее

Разница в 10+ раз - это не преувеличение. Python UDF создаёт Python worker-процесс, сериализует каждую строку через pickle, передаёт через IPC-сокет, выполняет Python-код, сериализует результат обратно. Higher-Order Function - это Java-код, который выполняется напрямую в JVM без каких-либо границ.

Практика: обработка заказов электронной коммерции

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

from pyspark.sql import SparkSession
from pyspark.sql import functions as F
from pyspark.sql.types import *

spark = SparkSession.builder.appName("ecommerce-complex-types").getOrCreate()

# Схема: каждый заказ содержит массив ItemType структур и словарь metadata
item_type = StructType([
    StructField("product_id",  StringType()),
    StructField("name",        StringType()),
    StructField("category",    StringType()),
    StructField("price",       DoubleType()),
    StructField("qty",         IntegerType()),
    StructField("in_stock",    BooleanType()),
])

order_schema = StructType([
    StructField("order_id",   StringType()),
    StructField("user_id",    IntegerType()),
    StructField("status",     StringType()),
    StructField("items",      ArrayType(item_type)),
    StructField("metadata",   MapType(StringType(), StringType())),
])

orders_raw = spark.createDataFrame([
    (
        "ORD-1001", 42, "confirmed",
        [
            ("P-100", "Ноутбук",      "electronics", 89999.0, 1, True),
            ("P-205", "Мышь",         "accessories", 1499.0,  2, True),
            ("P-301", "Сумка",        "accessories", 3990.0,  1, False),
        ],
        {"platform": "ios", "app_version": "3.4.1", "promo_code": "SALE20"},
    ),
    (
        "ORD-1002", 17, "pending",
        [
            ("P-400", "Наушники",     "electronics", 7990.0,  1, True),
            ("P-405", "Подставка",    "accessories", 2500.0,  1, True),
        ],
        {"platform": "web", "app_version": "n/a"},
    ),
    (
        "ORD-1003", 99, "confirmed",
        [
            ("P-600", "Планшет",      "electronics", 45000.0, 1, True),
            ("P-610", "Чехол",        "accessories", 890.0,   3, True),
            ("P-620", "Стилус",       "electronics", 4500.0,  2, True),
        ],
        {"platform": "android", "app_version": "3.4.0", "utm_source": "google"},
    ),
], order_schema)

# ── Задача 1: Оставить только доступные товары ────────────────────────────────
# filter() по полю in_stock внутри структуры
# items - это массив Struct, доступ к полю через .in_stock
available_orders = orders_raw.withColumn(
    "available_items",
    F.filter(
        F.col("items"),
        lambda item: item["in_stock"]   # обращение к полю структуры как к словарю
    )
)

available_orders.select(
    F.col("order_id"),
    F.size(F.col("items")).alias("total_items"),
    F.size(F.col("available_items")).alias("available_count"),
).show()
# +---------+-----------+---------------+
# |order_id |total_items|available_count|
# +---------+-----------+---------------+
# |ORD-1001 |3          |2              |   ← сумка не в наличии
# |ORD-1002 |2          |2              |
# |ORD-1003 |3          |3              |
# +---------+-----------+---------------+

# ── Задача 2: Сумма заказа со скидкой 20% для electronics ────────────────────
# aggregate() - проход по массиву, CASE для категорий
order_totals = available_orders.select(
    F.col("order_id"),
    F.aggregate(
        F.col("available_items"),
        F.lit(0.0),
        lambda acc, item: acc + (
            F.when(
                item["category"] == "electronics",
                F.round(item["price"] * item["qty"] * 0.8, 2)   # скидка 20%
            ).otherwise(
                item["price"] * item["qty"]                      # без скидки
            )
        )
    ).alias("total_with_discount")
)

order_totals.show()
# +---------+------------------+
# |order_id |total_with_discount|
# +---------+------------------+
# |ORD-1001 |74497.0            |   ← laptop×0.8 + mouse×2 + (bag отфильтрована)
# |ORD-1002 |8882.0             |   ← headphones×0.8 + stand
# |ORD-1003 |42960.0            |   ← tablet×0.8 + case×3 + stylus×2×0.8

# ── Задача 3: Привести названия к верхнему регистру ──────────────────────────
# transform() - модифицировать поле внутри каждого элемента структуры
# Нельзя просто upper(item) - нужно пересобрать структуру с изменённым полем

from pyspark.sql.functions import named_struct

orders_raw.select(
    F.col("order_id"),
    F.transform(
        F.col("items"),
        lambda item: named_struct(
            F.lit("product_id"),  item["product_id"],
            F.lit("name"),        F.upper(item["name"]),     # меняем только name
            F.lit("category"),    item["category"],
            F.lit("price"),       item["price"],
            F.lit("qty"),         item["qty"],
            F.lit("in_stock"),    item["in_stock"],
        )
    ).alias("items_upper_name")
).select(
    "order_id",
    F.col("items_upper_name")[0]["name"].alias("first_item_name")
).show()
# +---------+---------------+
# |order_id |first_item_name|
# +---------+---------------+
# |ORD-1001 |НОУТБУК        |
# |ORD-1002 |НАУШНИКИ       |
# |ORD-1003 |ПЛАНШЕТ        |
# +---------+---------+-----+

# ── Задача 4: Извлечь версию приложения из metadata ─────────────────────────
# Прямой доступ по ключу + coalesce для значения по умолчанию
from pyspark.sql.functions import coalesce

orders_raw.select(
    F.col("order_id"),
    F.col("metadata")["platform"].alias("platform"),
    coalesce(
        F.col("metadata")["app_version"],
        F.lit("unknown")
    ).alias("app_version"),
    F.col("metadata")["promo_code"].alias("promo_code"),  # None если нет ключа
).show()
# +---------+--------+-----------+----------+
# |order_id |platform|app_version|promo_code|
# +---------+--------+-----------+----------+
# |ORD-1001 |ios     |3.4.1      |SALE20    |
# |ORD-1002 |web     |n/a        |null      |
# |ORD-1003 |android |3.4.0      |null      |
# +---------+--------+-----------+----------+

# ── Задача 5: Есть ли в заказе хотя бы один electronics товар? ───────────────
orders_raw.select(
    F.col("order_id"),
    F.exists(
        F.col("items"),
        lambda item: item["category"] == "electronics"
    ).alias("has_electronics")
).show()
# +---------+---------------+
# |order_id |has_electronics|
# +---------+---------------+
# |ORD-1001 |true           |
# |ORD-1002 |true           |
# |ORD-1003 |true           |
# +---------+---------------+

explode vs Higher-Order Functions: когда что выбирать

Это один из самых важных практических вопросов при работе со сложными типами. Правило не абсолютное, но есть чёткие сигналы:

Операция Используй explode Используй HOF
JOIN по ключу элемента ❌ не подходит
Фильтрация элементов ❌ дорого filter()
Трансформация элементов ❌ дорого transform()
Сумма / агрегат по массиву ❌ дорого aggregate()
Проверка условия ❌ дорого exists() / forall()
Нужен индекс элемента ⚠️ posexplode ⚠️ 2-arg filter()
Конвертация в плоскую таблицу ✅ нужен результат ❌ не нужен

Best Practices

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

2. Предпочитайте StructType для известных схем, MapType для динамических. StructType поддерживает Nested Column Pruning в Parquet - чтение только нужных полей. MapType читается целиком. Если вы знаете набор ключей заранее - используйте StructType.

3. Higher-Order Functions вместо Python UDF для операций над массивами. Разница в производительности - 5–15× в пользу HOF. HOF компилируются в JVM-байткод Catalyst, UDF требует Python worker и pickle-сериализацию каждой строки.

4. collect_list + GroupBy вместо повторных join. Если вы агрегируете данные в витрину, соберите связанные данные в массив через collect_list, а не делайте join на каждый запрос. Это классическая денормализация Bronze→Silver.

5. explode_outer вместо explode если важны строки с пустыми массивами. explode молча удаляет строки с null или пустым массивом - это может быть неожиданной потерей данных. explode_outer сохраняет такие строки с null в колонке элемента.

6. Не читайте весь Map ради одного ключа - или переходите на StructType. Если вы всегда обращаетесь к metadata["version"] и metadata["platform"] - возможно, пора переводить эти поля в StructType с именованными полями. Это сэкономит I/O и позволит Parquet делать column pruning.