Complex Types: работа с ArrayType, MapType и функциями higher-order
ArrayType, MapType, StructType, explode, higher-order functions (transform, filter, aggregate) - обработка вложенных данных без взрыва строк
Зачем вообще нужны сложные типы¶
В классическом реляционном хранилище данные нормализованы: одна сущность - одна строка, одно значение - одна колонка. Это удобно для транзакций, но создаёт огромные трудности при работе с сырыми данными из реального мира.
Представьте типичный сценарий: пользователь делает заказ в интернет-магазине и покупает несколько товаров одновременно. В реляционной модели для этого нужны минимум три таблицы: 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.