Semi-structured данные: вложенные struct, flatten и нормализация
Kafka events, API payloads, вложенные JSON: как читать, анализировать, выравнивать и нормализовывать complex nested schema в production lakehouse pipeline.
Почему semi-structured данные стали стандартом¶
Современный Data Lake почти никогда не принимает аккуратные плоские таблицы. Реальные источники данных выглядят иначе:
- Kafka-события - вложенные JSON с метаданными, payload и заголовками
- REST API-ответы - объекты с несколькими уровнями вложенности, массивами и nullable-полями
- CDC-потоки (Change Data Capture) - события с before/after состоянием записи
- MongoDB/Firestore выгрузки - документы с произвольной структурой и вложенными объектами
- Webhook-payload - разные схемы в зависимости от типа события
Всё это semi-structured data: данные, у которых есть структура (поля с именами и типами), но структура неполная, изменчивая и иерархическая. Реляционная модель с фиксированными плоскими таблицами плохо подходит для хранения таких данных - отсюда рост JSON-based хранилищ и форматов.
Три типа данных и их место в lakehouse¶
Spark обрабатывает все три типа, но для semi-structured данных предоставляет особенно богатый инструментарий: StructType, ArrayType, MapType, функции from_json/to_json, explode, inline, flatten, higher-order functions.
Schema-on-read vs Schema-on-write¶
Schema-on-write (реляционные БД): схема определяется при создании таблицы. При записи несоответствующих данных - ошибка. Изменение схемы - ALTER TABLE.
Schema-on-read (Spark + Data Lake): данные хранятся как есть (JSON, Parquet, Avro). Схема применяется при чтении. Это даёт гибкость при приёме данных, но перекладывает ответственность за контроль качества на пайплайн.
Для semi-structured данных schema-on-read - это норма. Главное правило: чем раньше применить явную схему (StructType), тем лучше. Данные без схемы - это сырой JSON как один строковой тип; данные со схемой - полноценный DataFrame с типизированными колонками.
Типы данных для вложенных структур¶
Spark имеет три комплексных типа данных, которые позволяют описать любую вложенную структуру.
StructType - объект с именованными полями¶
StructType - это эквивалент JSON-объекта или Python-словаря с фиксированными ключами. Каждое поле имеет имя, тип и флаг nullable.
from pyspark.sql.types import *
# Схема пользователя с вложенным адресом
user_schema = StructType([
StructField("user_id", LongType(), nullable=False),
StructField("name", StringType()),
StructField("email", StringType()),
StructField("address", StructType([ # вложенный объект
StructField("city", StringType()),
StructField("country", StringType()),
StructField("zip", StringType()),
StructField("geo", StructType([ # вложенный внутри вложенного
StructField("lat", DoubleType()),
StructField("lon", DoubleType()),
])),
])),
StructField("created_at", TimestampType()),
])
В Parquet StructType хранится как группа колонок с иерархическими именами: address.city, address.geo.lat. Это фундаментальное отличие от реляционного хранения - связанные данные хранятся физически близко в одном файле.
ArrayType - упорядоченная коллекция¶
ArrayType содержит ноль или более элементов одного типа. Элементы могут быть как скалярными (ArrayType(StringType())), так и составными (ArrayType(StructType(...))).
StructField("tags", ArrayType(StringType())), # ['python', 'spark']
StructField("scores", ArrayType(DoubleType())), # [0.92, 0.87, 0.95]
StructField("orders", ArrayType(StructType([ # [{id: 1, amount: 500}, ...]
StructField("order_id", LongType()),
StructField("amount", DoubleType()),
]))),
ArrayType(StructType(...)) - самая частая комбинация в event-driven системах. Именно так хранятся позиции заказа, история событий, список адресов доставки.
MapType - словарь ключ-значение¶
MapType - пары ключ-значение, где ключи и значения имеют фиксированные типы.
StructField("metadata", MapType(StringType(), StringType())), # {"source": "web", "version": "2.1"}
StructField("counters", MapType(StringType(), LongType())), # {"clicks": 42, "views": 350}
MapType удобен для хранения разнородных атрибутов с переменным набором ключей - например, UTM-параметры, кастомные свойства событий, labels из Kubernetes.
Чтение и анализ вложенных данных¶
Чтение JSON с вложенной структурой¶
from pyspark.sql import SparkSession
from pyspark.sql.types import *
spark = SparkSession.builder.appName("nested-demo").getOrCreate()
# Типичный Kafka-payload: вложенный JSON
json_data = [
('{"event_id":"e001","type":"page_view","user":{"id":101,"segment":"premium"},"page":{"url":"/products","category":"electronics"},"device":{"os":"iOS","app_version":"3.2.1"},"ts":"2024-01-15T10:00:00Z"}',),
('{"event_id":"e002","type":"add_to_cart","user":{"id":102,"segment":"regular"},"page":{"url":"/cart","category":null},"device":{"os":"Android","app_version":"3.1.0"},"ts":"2024-01-15T10:01:00Z"}',),
('{"event_id":"e003","type":"purchase","user":{"id":101,"segment":"premium"},"page":{"url":"/checkout","category":"electronics"},"device":{"os":"iOS","app_version":"3.2.1"},"ts":"2024-01-15T10:05:00Z"}',),
]
raw_df = spark.createDataFrame(json_data, ["raw_json"])
Без схемы - это просто колонка строк. Применяем from_json с явной схемой:
from pyspark.sql.functions import from_json, col
event_schema = StructType([
StructField("event_id", StringType()),
StructField("type", StringType()),
StructField("user", StructType([
StructField("id", LongType()),
StructField("segment", StringType()),
])),
StructField("page", StructType([
StructField("url", StringType()),
StructField("category", StringType()),
])),
StructField("device", StructType([
StructField("os", StringType()),
StructField("app_version", StringType()),
])),
StructField("ts", StringType()),
])
events_df = raw_df.select(
from_json(col("raw_json"), event_schema).alias("event")
)
events_df.printSchema()
root
|-- event: struct (nullable = true)
| |-- event_id: string (nullable = true)
| |-- type: string (nullable = true)
| |-- user: struct (nullable = true)
| | |-- id: long (nullable = true)
| | |-- segment: string (nullable = true)
| |-- page: struct (nullable = true)
| | |-- url: string (nullable = true)
| | |-- category: string (nullable = true)
| |-- device: struct (nullable = true)
| | |-- os: string (nullable = true)
| | |-- app_version: string (nullable = true)
| |-- ts: string (nullable = true)
Почему явная схема важнее schema inference¶
Spark умеет автоматически определять схему: spark.read.json(path) или schema_of_json(). Но в production это антипаттерн:
Проблема 1: производительность. Schema inference читает всю выборку данных (по умолчанию 1000 строк, но при разнородных данных может больше) просто чтобы понять типы. Для больших файлов это дополнительный scan.
Проблема 2: тип по умолчанию. Если поле иногда "42" (строка), а иногда 42 (число), inference выбирает string. Явная схема позволяет принудительно задать LongType() - данные конвертируются при чтении.
Проблема 3: нестабильность. Сегодня inference даёт одну схему, завтра появятся новые события с новыми полями - схема изменится, и downstream-код сломается. Явная схема игнорирует неожиданные поля.
Проблема 4: NULL-handling. При inference все поля помечаются nullable. Явная схема позволяет указать nullable=False для обязательных полей и контролировать это.
printSchema и анализ структуры¶
printSchema() - первое что нужно сделать при работе с незнакомым датасетом:
# Полный вывод дерева схемы
df.printSchema()
# Программный обход схемы
def print_schema_recursive(schema, prefix=""):
for field in schema.fields:
field_type = field.dataType
nullable = "?" if field.nullable else "!"
print(f"{prefix}{field.name}: {field_type.typeName()}{nullable}")
if isinstance(field_type, StructType):
print_schema_recursive(field_type, prefix + " ")
elif isinstance(field_type, ArrayType) and isinstance(field_type.elementType, StructType):
print(f"{prefix} (array element):")
print_schema_recursive(field_type.elementType, prefix + " ")
print_schema_recursive(df.schema)
Навигация: точечная нотация и обращение к вложенным полям¶
Dot notation в DataFrame API¶
Доступ к вложенным полям - через точку, как в JavaScript:
# Извлекаем вложенные поля отдельно
flat_df = events_df.select(
col("event.event_id"),
col("event.type"),
col("event.user.id").alias("user_id"),
col("event.user.segment"),
col("event.page.url"),
col("event.page.category"),
col("event.device.os"),
col("event.device.app_version"),
col("event.ts"),
)
flat_df.show(truncate=False)
+--------+-----------+-------+---------+-----------+-----------+-------+-----------+-------------------+
|event_id| type|user_id| segment| url| category| os|app_version| ts|
+--------+-----------+-------+---------+-----------+-----------+-------+-----------+-------------------+
| e001| page_view| 101| premium| /products|electronics| iOS| 3.2.1|2024-01-15T10:00:00Z|
| e002|add_to_cart| 102| regular| /cart| null|Android| 3.1.0|2024-01-15T10:01:00Z|
| e003| purchase| 101| premium| /checkout|electronics| iOS| 3.2.1|2024-01-15T10:05:00Z|
+--------+-----------+-------+---------+-----------+-----------+-------+-----------+-------------------+
Точечная нотация работает на произвольной глубине: col("event.user.address.geo.lat") - валидное выражение.
Dot notation в Spark SQL¶
В SQL всё работает так же интуитивно:
-- Регистрируем как временное представление
events_df.createOrReplaceTempView("events")
SELECT
event.event_id,
event.user.id AS user_id,
event.user.segment,
event.device.os,
event.page.url
FROM events
WHERE event.user.segment = 'premium'
AND event.type = 'purchase'
SQL-синтаксис через точку работает для StructType на любой глубине. Для MapType используется map_column['key'] или map_column.key (если ключи - идентификаторы).
Доступ к полям MapType¶
from pyspark.sql.functions import col, map_keys, map_values
# Предположим, metadata - это MapType(StringType, StringType)
df.select(
col("metadata")["source"].alias("source"), # по ключу
col("metadata")["utm_campaign"].alias("campaign"),
map_keys("metadata").alias("all_keys"), # все ключи как array
map_values("metadata").alias("all_values"), # все значения как array
)
Работа с ArrayType: базовые операции без explode¶
До разворачивания массива есть много полезных операций:
from pyspark.sql.functions import size, array_contains, element_at, slice as spark_slice
df.select(
"user_id",
size("tags").alias("tags_count"), # длина массива
array_contains("tags", "premium").alias("is_premium"), # содержит ли элемент
element_at("tags", 1).alias("first_tag"), # первый элемент (1-indexed!)
spark_slice("scores", 1, 3).alias("top3_scores"), # срез массива
)
element_at использует 1-based индексацию (как SQL), а не 0-based (как Python). Отрицательные индексы работают: element_at("tags", -1) - последний элемент.
Как Parquet хранит вложенные данные¶
Понимание физического хранения помогает писать оптимальные запросы.
Dremel-style columnar encoding¶
Parquet использует алгоритм из статьи Google Dremel для хранения вложенных данных. Каждый "листовой" атрибут вложенной структуры хранится как отдельная колонка, плюс два служебных массива:
- Definition levels: для каждого значения - как глубоко в иерархии оно определено (чтобы отличить null от "отсутствующего родителя")
- Repetition levels: для массивов - признак повторения на каком уровне
Для user.address.geo.lat (4 уровня) Parquet хранит:
Column: user.address.geo.lat
Definition levels: [3, 3, 2, 1, 3, ...]
↑ ↑ ↑ ↑
lat определён | адрес есть, но geo = null | адрес null | user null
Это позволяет Parquet читать только нужную колонку без чтения всей записи. При SELECT user.address.city FROM ... Spark читает только один столбец, игнорируя user.address.geo.lat, user.name и все остальные.
Column pruning для nested fields¶
Catalyst-оптимизатор Spark умеет делать column pruning (отсечение ненужных колонок) даже для вложенных структур:
# Spark прочитает из Parquet ТОЛЬКО колонки user.id и user.segment
# Остальные колонки (page, device, ts) не читаются с диска
events_df.select("event.user.id", "event.user.segment").explain()
В плане выполнения увидите:
FileScan parquet [...] PushedFilters: [...], ReadSchema: struct<user:struct<id:bigint,segment:string>>
ReadSchema показывает что именно читается. Это ключевое преимущество вложенных структур в Parquet: локальность данных + column pruning = минимальный I/O.
Predicate pushdown для nested fields¶
Предикаты на вложенных полях тоже могут быть pushed down:
events_df.filter(col("event.user.segment") == "premium").explain()
FileScan parquet [...] PushedFilters: [IsNotNull(event.user.segment), EqualTo(event.user.segment, premium)]
Parquet-файл будет отфильтрован на уровне row group statistics - если весь row group не содержит premium, он пропускается без чтения.
Flatten: выравнивание вложенных структур¶
Зачем нужен flatten¶
Вложенные структуры отлично подходят для хранения и Spark-обработки. Но аналитики и BI-инструменты часто ожидают плоские таблицы:
- SQL-аналитики привыкают к
SELECT city FROM users, а неSELECT address.city FROM users - BI-инструменты (Tableau, Power BI, Metabase) часто не могут работать с вложенными полями напрямую
- ML-фреймворки ожидают flat feature matrix - один вектор строка с числовыми колонками
- JOIN-операции проще на плоских таблицах
Flatten - это преобразование иерархической схемы в плоскую, где каждый "листовой" атрибут становится колонкой верхнего уровня.
Ручной flatten - антипаттерн¶
Первое, что приходит в голову - прописать все поля вручную:
# АНТИПАТТЕРН: хрупко, не масштабируется
flat = df.select(
col("event.event_id"),
col("event.type"),
col("event.user.id").alias("user_id"),
col("event.user.segment").alias("user_segment"),
col("event.page.url").alias("page_url"),
col("event.page.category").alias("page_category"),
col("event.device.os").alias("device_os"),
col("event.device.app_version").alias("device_app_version"),
col("event.ts"),
)
Проблемы этого подхода:
- При добавлении нового поля в схему источника (например,
device.model) - нужно вручную редактировать пайплайн - При 4+ уровнях вложенности список алиасов становится неподъёмным
- Нет стандарта именования: один разработчик напишет
user_id, другой -userId, третий -usr_id
Универсальный рекурсивный flatten¶
Правильное решение - написать функцию, которая обходит схему и автоматически генерирует select-выражения:
from pyspark.sql import DataFrame
from pyspark.sql.types import StructType, ArrayType
from pyspark.sql.functions import col
def flatten_struct(schema: StructType, prefix: str = "", sep: str = "_") -> list:
"""
Рекурсивно обходит StructType и возвращает список Column-выражений
с выровненными именами. ArrayType и MapType оставляются как есть —
их нельзя выровнять без explode/unnesting.
"""
columns = []
for field in schema.fields:
full_path = f"{prefix}{sep}{field.name}" if prefix else field.name
field_col = col(f"`{full_path}`") if not prefix else col(f"`{prefix}`.`{field.name}`")
if isinstance(field.dataType, StructType):
# Рекурсивно идём вглубь структуры
nested = flatten_struct(field.dataType, full_path, sep)
columns.extend(nested)
else:
# Листовое поле: берём как есть, даём "плоское" имя
# Заменяем точки на разделитель для безопасного имени колонки
safe_name = full_path.replace(".", sep)
columns.append(col(full_path).alias(safe_name))
return columns
def flatten_df(df: DataFrame, sep: str = "_") -> DataFrame:
"""
Выравнивает все StructType-колонки датафрейма до одного уровня.
ArrayType и MapType остаются нетронутыми.
"""
flat_columns = []
for field in df.schema.fields:
if isinstance(field.dataType, StructType):
nested = flatten_struct(field.dataType, field.name, sep)
flat_columns.extend(nested)
else:
flat_columns.append(col(field.name))
return df.select(flat_columns)
Применяем:
# events_df имеет одну колонку "event" типа StructType
flat_events = flatten_df(events_df)
flat_events.printSchema()
root
|-- event_event_id: string (nullable = true)
|-- event_type: string (nullable = true)
|-- event_user_id: long (nullable = true)
|-- event_user_segment: string (nullable = true)
|-- event_page_url: string (nullable = true)
|-- event_page_category: string (nullable = true)
|-- event_device_os: string (nullable = true)
|-- event_device_app_version: string (nullable = true)
|-- event_ts: string (nullable = true)
Всё дерево вложенности превратилось в плоские колонки с именами через разделитель _. Никакого хардкода полей.
Улучшенная версия с snake_case и trim prefix¶
На практике хочется убрать лишний префикс верхнего уровня и привести имена к snake_case:
import re
def to_snake_case(name: str) -> str:
"""CamelCase и PascalCase → snake_case"""
s1 = re.sub(r'(.)([A-Z][a-z]+)', r'\1_\2', name)
return re.sub(r'([a-z0-9])([A-Z])', r'\1_\2', s1).lower()
def flatten_struct_v2(schema: StructType, prefix: str = "", sep: str = "_",
strip_prefix: str = "") -> list:
"""
Расширенная версия: поддерживает strip_prefix для удаления лишнего
корневого уровня и to_snake_case для нормализации имён.
"""
columns = []
for field in schema.fields:
full_path = f"{prefix}.{field.name}" if prefix else field.name
if isinstance(field.dataType, StructType):
nested = flatten_struct_v2(field.dataType, full_path, sep, strip_prefix)
columns.extend(nested)
else:
# Формируем "чистое" имя: убираем strip_prefix, меняем точки на _
clean_path = full_path
if strip_prefix and clean_path.startswith(strip_prefix + "."):
clean_path = clean_path[len(strip_prefix) + 1:]
col_name = to_snake_case(clean_path.replace(".", sep))
columns.append(col(full_path).alias(col_name))
return columns
def flatten_df_v2(df: DataFrame, sep: str = "_", strip_top_prefix: bool = True) -> DataFrame:
flat_columns = []
for field in df.schema.fields:
if isinstance(field.dataType, StructType):
# strip_prefix = имя корневой колонки (например, "event")
prefix = field.name if strip_top_prefix else ""
nested = flatten_struct_v2(field.dataType, field.name, sep, strip_prefix=prefix)
flat_columns.extend(nested)
else:
flat_columns.append(col(field.name))
return df.select(flat_columns)
flat_clean = flatten_df_v2(events_df, strip_top_prefix=True)
flat_clean.printSchema()
root
|-- event_id: string
|-- type: string
|-- user_id: long
|-- user_segment: string
|-- page_url: string
|-- page_category: string
|-- device_os: string
|-- device_app_version: string
|-- ts: string
Префикс event_ убран (так как мы читаем из поля event), имена в нижнем регистре.
Ограничения flatten: ArrayType¶
Важное ограничение: функция flatten_df не может автоматически обработать ArrayType. Массив нельзя "выровнять" - у него нет фиксированного числа элементов, нельзя создать items_0, items_1, items_2... колонки (их количество у каждой строки разное).
# Если в схеме есть ArrayType - он остаётся как есть
schema_with_array = StructType([
StructField("user_id", LongType()),
StructField("profile", StructType([
StructField("name", StringType()),
StructField("age", IntegerType()),
])),
StructField("tags", ArrayType(StringType())), # ← остаётся массивом
])
df = spark.createDataFrame(
[(1, ("Alice", 30), ["python", "spark"]),
(2, ("Bob", 25), ["java"])],
schema_with_array
)
flatten_df(df).printSchema()
root
|-- user_id: long
|-- profile_name: string
|-- profile_age: integer
|-- tags: array<string> ← НЕ развёрнут
Для обработки ArrayType нужен отдельный шаг: explode (если нужны строки) или higher-order functions (если нужна трансформация внутри массива без разворачивания).
Multi-level flattening: глубоко вложенные структуры¶
Пример 4-уровневой вложенности¶
deep_schema = StructType([
StructField("order_id", LongType()),
StructField("customer", StructType([
StructField("id", LongType()),
StructField("name", StringType()),
StructField("contact", StructType([
StructField("email", StringType()),
StructField("phone", StructType([
StructField("country_code", StringType()),
StructField("number", StringType()),
])),
])),
])),
StructField("shipping", StructType([
StructField("address", StructType([
StructField("street", StringType()),
StructField("city", StringType()),
StructField("country", StringType()),
])),
StructField("method", StringType()),
StructField("estimated_days", IntegerType()),
])),
])
data = [(
1001,
(42, "Alice", ("alice@example.com", ("+7", "9161234567"))),
(("Ленина 1", "Москва", "RU"), "express", 2),
)]
deep_df = spark.createDataFrame(data, deep_schema)
deep_df.printSchema()
root
|-- order_id: long (nullable = true)
|-- customer: struct (nullable = true)
| |-- id: long (nullable = true)
| |-- name: string (nullable = true)
| |-- contact: struct (nullable = true)
| | |-- email: string (nullable = true)
| | |-- phone: struct (nullable = true)
| | | |-- country_code: string (nullable = true)
| | | |-- number: string (nullable = true)
|-- shipping: struct (nullable = true)
| |-- address: struct (nullable = true)
| | |-- street: string (nullable = true)
| | |-- city: string (nullable = true)
| | |-- country: string (nullable = true)
| |-- method: string (nullable = true)
| |-- estimated_days: integer (nullable = true)
Применяем нашу рекурсивную функцию:
flat_deep = flatten_df(deep_df)
flat_deep.printSchema()
root
|-- order_id: long
|-- customer_id: long
|-- customer_name: string
|-- customer_contact_email: string
|-- customer_contact_phone_country_code: string
|-- customer_contact_phone_number: string
|-- shipping_address_street: string
|-- shipping_address_city: string
|-- shipping_address_country: string
|-- shipping_method: string
|-- shipping_estimated_days: integer
Все 4 уровня вложенности развёрнуты автоматически. Код пайплайна не содержит ни одного захардкоженного имени поля.
Нормализация: стратегии проектирования Silver-слоя¶
Flatten - это не единственный способ работы с вложенными данными. Иногда правильнее нормализовать данные в несколько связанных таблиц, а иногда - сохранить часть вложенности и разложить только нужное.
Когда сохранять вложенность¶
Вложенная структура полезна когда связанные данные читаются вместе. Если аналитика всегда читает customer.name и customer.email вместе - нет смысла разносить их в разные колонки верхнего уровня.
# Selective flatten: раскрываем только то, что нужно аналитике
# customer и device остаются как structs - они читаются как единица
silver_events = events_df.select(
col("event.event_id"),
col("event.type").alias("event_type"),
col("event.user").alias("user"), # struct сохраняется как есть
col("event.page.url").alias("page_url"),
col("event.page.category").alias("page_category"),
col("event.device").alias("device"), # struct сохраняется
col("event.ts").cast("timestamp").alias("event_ts"),
)
silver_events.printSchema()
root
|-- event_id: string
|-- event_type: string
|-- user: struct
| |-- id: long
| |-- segment: string
|-- page_url: string
|-- page_category: string
|-- device: struct
| |-- os: string
| |-- app_version: string
|-- event_ts: timestamp
Такая схема - хороший Silver. Аналитик может писать WHERE user.segment = 'premium' - понятно и удобно. И Parquet будет читать только нужные колонки.
Нормализация в отдельные таблицы¶
При глубокой вложенности с разными сущностями правильно создать отдельные таблицы:
# У нас есть deep_df с customer, shipping, items (array)
# Нормализуем в три таблицы Silver
# Silver: заказы (шапка)
silver_orders = deep_df.select(
col("order_id"),
col("customer.id").alias("customer_id"),
col("shipping.method").alias("shipping_method"),
col("shipping.estimated_days"),
)
# Silver: клиенты (deduplicated)
silver_customers = deep_df.select(
col("customer.id").alias("customer_id"),
col("customer.name").alias("customer_name"),
col("customer.contact.email").alias("email"),
).dropDuplicates(["customer_id"])
# Silver: адреса доставки
silver_shipping_addresses = deep_df.select(
col("order_id"),
col("shipping.address.street"),
col("shipping.address.city"),
col("shipping.address.country"),
)
Денормализация: struct() для обратной сборки¶
Иногда нужно сделать обратное: взять набор плоских колонок и упаковать их в struct. Это денормализация - полезна для Gold-слоя, где данные компактнее хранить в logical groups.
from pyspark.sql.functions import struct, col
# Пришли плоские данные из разных источников
flat_df = spark.createDataFrame([
(1, "Alice", "alice@example.com", "+7", "9161234567", "Moscow", "RU"),
], ["customer_id", "name", "email", "phone_code", "phone_num", "city", "country"])
# Собираем логические блоки через struct()
denormalized = flat_df.select(
col("customer_id"),
struct(
col("name"),
col("email"),
struct(
col("phone_code").alias("country_code"),
col("phone_num").alias("number"),
).alias("phone"),
).alias("customer"),
struct(
col("city"),
col("country"),
).alias("location"),
)
denormalized.printSchema()
root
|-- customer_id: long
|-- customer: struct
| |-- name: string
| |-- email: string
| |-- phone: struct
| | |-- country_code: string
| | |-- number: string
|-- location: struct
| |-- city: string
| |-- country: string
|-- ts: timestamp
struct() - это сборка нескольких колонок в один составной тип. Полезно когда:
- Gold-слой должен содержать меньше колонок верхнего уровня
- Нужно передать данные в систему, ожидающую JSON-подобные объекты
- Группируете связанные атрибуты для ясности схемы
Практика: полный пайплайн Bronze → Silver¶
Разберём реалистичный кейс: выгрузка пользователей из MongoDB.
Исходные данные (имитация MongoDB export)¶
from pyspark.sql.types import *
from pyspark.sql.functions import *
# Полная схема документа из MongoDB
mongo_schema = StructType([
StructField("_id", StringType()),
StructField("userId", LongType()),
StructField("createdAt", StringType()),
StructField("updatedAt", StringType()),
StructField("profile", StructType([
StructField("firstName", StringType()),
StructField("lastName", StringType()),
StructField("birthDate", StringType()),
StructField("language", StringType()),
StructField("avatarUrl", StringType()),
])),
StructField("subscription", StructType([
StructField("plan", StringType()),
StructField("startedAt", StringType()),
StructField("expiresAt", StringType()),
StructField("autoRenew", BooleanType()),
StructField("features", ArrayType(StringType())),
])),
StructField("device", StructType([
StructField("type", StringType()),
StructField("os", StringType()),
StructField("appVersion", StringType()),
StructField("pushEnabled", BooleanType()),
])),
StructField("geo", StructType([
StructField("country", StringType()),
StructField("city", StringType()),
StructField("lat", DoubleType()),
StructField("lon", DoubleType()),
StructField("timezone", StringType()),
])),
StructField("stats", StructType([
StructField("sessionsTotal", LongType()),
StructField("lastLoginAt", StringType()),
StructField("purchasesCount", IntegerType()),
StructField("totalSpent", DoubleType()),
])),
])
# Тестовые данные
users_data = [
("507f1f77bcf86cd799439011", 1001, "2023-01-10T08:00:00Z", "2024-01-15T12:30:00Z",
("Alice", "Smith", "1990-05-15", "ru", "https://cdn.example.com/avatars/1001.jpg"),
("premium", "2023-06-01T00:00:00Z", "2024-06-01T00:00:00Z", True, ["analytics", "export", "api"]),
("mobile", "iOS", "3.2.1", True),
("RU", "Moscow", 55.7558, 37.6173, "Europe/Moscow"),
(245, "2024-01-15T10:00:00Z", 12, 15670.50)),
("507f1f77bcf86cd799439012", 1002, "2023-03-20T09:00:00Z", "2024-01-14T08:00:00Z",
("Bob", "Johnson", "1985-11-22", "en", None),
("free", "2023-03-20T00:00:00Z", None, False, ["analytics"]),
("desktop", "Windows", "3.1.0", False),
("US", "New York", 40.7128, -74.0060, "America/New_York"),
(42, "2024-01-14T08:00:00Z", 0, 0.0)),
]
bronze_users = spark.createDataFrame(users_data, mongo_schema)
Шаг 1: Анализ входящей схемы¶
# Всегда начинаем с анализа
bronze_users.printSchema()
bronze_users.show(2, truncate=False, vertical=True)
# Проверяем nullability и типы
bronze_users.select(
[sum(col(c).isNull().cast("int")).alias(c) for c in bronze_users.columns]
).show()
Шаг 2: Silver - нормализованная таблица пользователей¶
from pyspark.sql.functions import col, to_timestamp, concat_ws, lower, trim
silver_users = bronze_users.select(
# Основные поля
col("userId").alias("user_id"),
to_timestamp(col("createdAt"), "yyyy-MM-dd'T'HH:mm:ss'Z'").alias("created_at"),
to_timestamp(col("updatedAt"), "yyyy-MM-dd'T'HH:mm:ss'Z'").alias("updated_at"),
# Профиль: раскрываем, нормализуем имена (camelCase → snake_case)
concat_ws(" ", col("profile.firstName"), col("profile.lastName")).alias("full_name"),
col("profile.firstName").alias("first_name"),
col("profile.lastName").alias("last_name"),
col("profile.birthDate").alias("birth_date"),
lower(trim(col("profile.language"))).alias("language"),
# Подписка: раскрываем ключевые поля
col("subscription.plan").alias("subscription_plan"),
to_timestamp(col("subscription.startedAt"), "yyyy-MM-dd'T'HH:mm:ss'Z'").alias("subscription_started_at"),
to_timestamp(col("subscription.expiresAt"), "yyyy-MM-dd'T'HH:mm:ss'Z'").alias("subscription_expires_at"),
col("subscription.autoRenew").alias("auto_renew"),
col("subscription.features").alias("subscription_features"), # array остаётся
# Device: оставляем как struct для экономии колонок
col("device"),
# Geo: раскрываем страну/город, координаты в struct
col("geo.country").alias("country"),
col("geo.city").alias("city"),
col("geo.timezone").alias("timezone"),
struct(
col("geo.lat").alias("lat"),
col("geo.lon").alias("lon"),
).alias("coordinates"),
# Stats
col("stats.sessionsTotal").alias("sessions_total"),
to_timestamp(col("stats.lastLoginAt"), "yyyy-MM-dd'T'HH:mm:ss'Z'").alias("last_login_at"),
col("stats.purchasesCount").alias("purchases_count"),
col("stats.totalSpent").alias("total_spent"),
)
silver_users.printSchema()
root
|-- user_id: long
|-- created_at: timestamp
|-- updated_at: timestamp
|-- full_name: string
|-- first_name: string
|-- last_name: string
|-- birth_date: string
|-- language: string
|-- subscription_plan: string
|-- subscription_started_at: timestamp
|-- subscription_expires_at: timestamp
|-- auto_renew: boolean
|-- subscription_features: array<string>
|-- device: struct
| |-- type: string
| |-- os: string
| |-- appVersion: string
| |-- pushEnabled: boolean
|-- country: string
|-- city: string
|-- timezone: string
|-- coordinates: struct
| |-- lat: double
| |-- lon: double
|-- sessions_total: long
|-- last_login_at: timestamp
|-- purchases_count: integer
|-- total_spent: double
Это хорошая Silver-схема: основные поля раскрыты и нормализованы, device и coordinates сохранены как structs (они читаются вместе), subscription_features остался массивом.
Шаг 3: Gold - аналитическая витрина¶
from pyspark.sql.functions import datediff, current_date, array_size, when
gold_user_segments = silver_users.select(
col("user_id"),
col("full_name"),
col("country"),
col("subscription_plan"),
# Активность
datediff(current_date(), col("last_login_at")).alias("days_since_login"),
col("sessions_total"),
col("purchases_count"),
col("total_spent"),
# Признаки для сегментации
array_size("subscription_features").alias("feature_count"),
col("device.os").alias("os"),
col("device.type").alias("device_type"),
# Сегмент
when(col("subscription_plan") == "premium", "premium_user")
.when(col("purchases_count") > 5, "active_free")
.otherwise("passive_free")
.alias("user_segment"),
)
gold_user_segments.show(truncate=False)
Шаг 4: Unnesting массива фич¶
Для аналитики по фичам подписки нужен explode:
from pyspark.sql.functions import explode
# Какие фичи наиболее популярны?
silver_users.select(
"user_id",
"subscription_plan",
explode("subscription_features").alias("feature")
).groupBy("feature", "subscription_plan") \
.count() \
.orderBy("count", ascending=False) \
.show()
Nullable semantics в вложенных схемах¶
Почти всё в вложенных схемах - nullable. Это важно понимать:
Два вида null для вложенных полей¶
# Null 1: struct сам по себе null
data = [(1, None)] # весь geo = null
df = spark.createDataFrame(data, StructType([
StructField("user_id", LongType()),
StructField("geo", StructType([
StructField("city", StringType()),
StructField("country", StringType()),
])),
]))
# При обращении к geo.city когда geo = null → тоже null, без ошибки
df.select("user_id", col("geo.city").alias("city")).show()
# +-------+----+
# |user_id|city|
# +-------+----+
# | 1|null|
# +-------+----+
# Null 2: struct существует, но поле внутри null
data2 = [(2, ("Moscow", None))] # geo.country = null
df2 = spark.createDataFrame(data2, StructType([
StructField("user_id", LongType()),
StructField("geo", StructType([
StructField("city", StringType()),
StructField("country", StringType()),
])),
]))
Spark корректно обрабатывает оба случая: обращение к полю через null-родителя возвращает null, а не ошибку. Это удобно, но важно помнить при фильтрации:
# Это вернёт строки где geo.city = null ИЛИ geo = null
df.filter(col("geo.city").isNull())
# Если хотите только строки где geo существует, но city = null:
df.filter(col("geo").isNotNull() & col("geo.city").isNull())
coalesce для safe extraction¶
from pyspark.sql.functions import coalesce, lit
# Безопасное извлечение: fallback на "unknown"
df.select(
"user_id",
coalesce(col("geo.city"), lit("unknown")).alias("city"),
coalesce(col("geo.country"), lit("XX")).alias("country_code"),
)
Schema evolution при работе со struct¶
Delta Lake поддерживает schema evolution для вложенных структур, но с ограничениями.
Добавление новых полей в struct¶
# Запись первой версии схемы
users_v1 = spark.createDataFrame(
[(1, ("Alice", "Smith"))],
StructType([
StructField("user_id", LongType()),
StructField("profile", StructType([
StructField("first_name", StringType()),
StructField("last_name", StringType()),
])),
])
)
users_v1.write.format("delta").save("/tmp/delta/users")
# Новая версия с дополнительным полем в profile
users_v2 = spark.createDataFrame(
[(2, ("Bob", "Johnson", "bob@example.com"))],
StructType([
StructField("user_id", LongType()),
StructField("profile", StructType([
StructField("first_name", StringType()),
StructField("last_name", StringType()),
StructField("email", StringType()), # новое поле
])),
])
)
# mergeSchema позволяет добавить новое поле
users_v2.write.format("delta") \
.mode("append") \
.option("mergeSchema", "true") \
.save("/tmp/delta/users")
# Теперь при чтении старых записей email = null
spark.read.format("delta").load("/tmp/delta/users").show()
+-------+--------------------+
|user_id| profile|
+-------+--------------------+
| 1|{Alice, Smith, null}| ← email = null для старой записи
| 2|{Bob, Johnson, bob@...|
+-------+--------------------+
Ограничения schema evolution для nested struct¶
Delta Lake поддерживает добавление новых полей в struct через mergeSchema. Но:
- Удаление полей не поддерживается автоматически (нужен
overwriteSchema) - Изменение типа поля в struct - только для совместимых типов (int → long)
- Переименование поля не поддерживается как evolution - только через полную перезапись
Performance: Catalyst и nested data¶
Nested field pruning в EXPLAIN¶
silver_users.select("user_id", "country", "device.os").explain(mode="formatted")
В плане выполнения ищите ReadSchema:
FileScan parquet [...] ReadSchema: struct<user_id:bigint, country:string, device:struct<os:string>>
Spark читает из Parquet только три "пути": user_id, country, device.os. Все остальные колонки (profile, subscription, geo, stats) физически не читаются.
Когда column pruning не работает¶
Column pruning ломается если вы сначала читаете весь struct, а потом фильтруете:
# ПЛОХО: читается весь device-struct, потом берётся только os
silver_users.select("user_id", col("device").getField("os"))
# ХОРОШО: pruning работает корректно
silver_users.select("user_id", col("device.os"))
На уровне логического плана результат одинаков, но первый вариант в некоторых версиях Spark может не оптимизировать pruning правильно. Используйте dot-notation - это явный сигнал Catalyst что нужно только это поле.
Кэширование вложенных датафреймов¶
# Если используете nested df в нескольких местах - кэшируйте
silver_users.cache()
silver_users.count() # trigger materialization
# Теперь все последующие запросы к silver_users читают из памяти
silver_users.filter(col("subscription_plan") == "premium").count()
silver_users.groupBy("country").count().show()
silver_users.unpersist()
Anti-patterns¶
Anti-pattern 1: Raw JSON в Silver и Gold¶
# НЕПРАВИЛЬНО: хранить raw JSON строку в Silver
silver.withColumn("raw_payload", to_json(col("event"))) # зачем?
# В Silver хранятся структурированные данные
# Raw JSON - только в Bronze, как fallback или аудит-трейл
Anti-pattern 2: Повторный explode при каждом запросе¶
# НЕПРАВИЛЬНО: explode в ad-hoc запросах на raw Bronze
bronze_orders \
.select("order_id", explode("items")) \
.filter(col("col.category") == "electronics") \
.show()
# ПРАВИЛЬНО: explode один раз в Silver, хранить нормализованную таблицу
# silver_order_items уже нормализована - делаем обычный filter
silver_order_items.filter(col("category") == "electronics").show()
Anti-pattern 3: Giant wide flattened table на Bronze¶
# НЕПРАВИЛЬНО: делать полный flatten сразу при ingestion в Bronze
raw_events.transform(flatten_df).write.format("delta").save("bronze/events")
# Bronze должен хранить данные максимально близко к исходному формату
# Flatten - операция Silver или Gold уровня
Anti-pattern 4: Схема без explicit typing¶
# НЕПРАВИЛЬНО: schema inference для production pipeline
df = spark.read.json("s3a://bucket/events/")
# Каждый запуск может дать разную схему, типы нестабильны
# ПРАВИЛЬНО: всегда явная схема
df = spark.read.schema(event_schema).json("s3a://bucket/events/")
Anti-pattern 5: Хардкод имён nested полей¶
# НЕПРАВИЛЬНО: список полей захардкожен
flat = df.select(
col("profile.firstName"),
col("profile.lastName"),
# ... 40 строк
)
# ПРАВИЛЬНО: рекурсивный flatten или генерация из schema
flat = flatten_df(df)
# или
fields = get_nested_fields(df.schema, "profile")
flat = df.select(*fields)
Домашнее задание¶
Датасет: глубоко вложенный JSON из IoT-платформы:
iot_schema = StructType([
StructField("deviceId", StringType()),
StructField("timestamp", StringType()),
StructField("location", StructType([
StructField("plant", StringType()),
StructField("zone", StringType()),
StructField("coordinates", StructType([
StructField("x", DoubleType()),
StructField("y", DoubleType()),
StructField("floor", IntegerType()),
])),
])),
StructField("sensors", StructType([
StructField("environment", StructType([
StructField("tempC", DoubleType()),
StructField("humidityPct", DoubleType()),
StructField("pressureHpa", DoubleType()),
])),
StructField("mechanical", StructType([
StructField("vibrationHz", DoubleType()),
StructField("rotationRpm", IntegerType()),
StructField("loadPct", DoubleType()),
])),
])),
StructField("alerts", ArrayType(StructType([
StructField("code", StringType()),
StructField("severity", StringType()),
StructField("message", StringType()),
]))),
StructField("meta", MapType(StringType(), StringType())),
])
Задание:
- Написать функцию
auto_flatten(df), которая рекурсивно выравнивает всеStructType-колонки, при этом: - Переименовывает все колонки в
snake_case - Убирает лишние префиксы (
sensors_environment_temp_c→temp_cилиenv_temp_c) -
ArrayTypeиMapTypeоставляет как есть -
Применить
auto_flattenк датасету - в коде не должно быть ни одного захардкоженного имени поля -
Создать три аналитические агрегации поверх выровненного датасета:
- Средняя температура по
location.zoneза последний час - Процент устройств с алертами severity=
criticalпоlocation.plant -
Top-5 устройств по максимальному значению
vibration_hz -
Из массива
alertsсоздать отдельную нормализованную таблицуdevice_alerts(device_id, alert_code, severity, message, ts)черезexplode -
Для каждого устройства собрать обратно struct
device_healthс полямиtemp_c,vibration_hz,load_pctчерезstruct()- сохранить в Gold-таблицу
Чеклист¶
- [ ] Понимаю разницу между StructType, ArrayType и MapType - когда каждый тип уместен
- [ ] Умею задавать явную схему через
StructType(не использую schema inference в production) - [ ] Читаю и понимаю
printSchema()для глубоко вложенных датасетов - [ ] Знаю dot notation:
col("parent.child.field")для StructType иmap_col["key"]для MapType - [ ] Знаю как Parquet хранит nested данные (Dremel encoding) и почему column pruning работает
- [ ] Написал / понимаю рекурсивную функцию
flatten_df- без хардкода имён полей - [ ] Знаю ограничение flatten: ArrayType нельзя выровнять без explode
- [ ] Понимаю разницу между flatten, нормализацией и денормализацией - когда что применять
- [ ] Умею собрать struct обратно через
struct()для денормализации - [ ] Знаю nullable semantics: null-struct vs null-field, безопасный доступ через точку
- [ ] Понимаю schema evolution в Delta Lake для вложенных структур (mergeSchema)
- [ ] Не делаю повторный explode в аналитике - нормализую один раз в Silver