Semi-structured данные: вложенные struct, flatten и нормализация

Kafka events, API payloads, вложенные JSON: как читать, анализировать, выравнивать и нормализовывать complex nested schema в production lakehouse pipeline.

core

Почему 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"),
)

Проблемы этого подхода:

  1. При добавлении нового поля в схему источника (например, device.model) - нужно вручную редактировать пайплайн
  2. При 4+ уровнях вложенности список алиасов становится неподъёмным
  3. Нет стандарта именования: один разработчик напишет 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())),
])

Задание:

  1. Написать функцию auto_flatten(df), которая рекурсивно выравнивает все StructType-колонки, при этом:
  2. Переименовывает все колонки в snake_case
  3. Убирает лишние префиксы (sensors_environment_temp_ctemp_c или env_temp_c)
  4. ArrayType и MapType оставляет как есть

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

  6. Создать три аналитические агрегации поверх выровненного датасета:

  7. Средняя температура по location.zone за последний час
  8. Процент устройств с алертами severity=critical по location.plant
  9. Top-5 устройств по максимальному значению vibration_hz

  10. Из массива alerts создать отдельную нормализованную таблицу device_alerts(device_id, alert_code, severity, message, ts) через explode

  11. Для каждого устройства собрать обратно 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