JSON: from_json, to_json, get_json_object, schema_of_json

from_json и schema_of_json для парсинга, get_json_object vs json_tuple, to_json для генерации, обработка битых записей, хранение JSON в Parquet и ограничения Predicate Pushdown.

core

JSON в мире Big Data: архитектурный вызов

JSON (JavaScript Object Notation) - универсальный язык обмена данными в modern data stack. REST API возвращают JSON. Kafka-события - JSON. CDC-стримы из баз данных через Debezium - JSON. Мобильные приложения отправляют JSON-события в аналитический pipeline. Логи веб-серверов в формате JSONL. Webhooks - JSON.

Проблема в том, что JSON - это текст. Он человекочитаем и гибок, но аналитические движки, включая Spark, работают с ним неэффективно:

CPU-bound парсинг. Каждый раз когда нужно извлечь поле из JSON-строки, Spark запускает лексер/парсер Jackson (Java-библиотека для JSON). Для одного поля это быстро. Для миллиарда строк - дорого. В отличие от Parquet, где значение поля лежит в известном байтовом смещении в файле, JSON-поле нужно искать текстовым разбором.

Нет колоночного сжатия. Parquet хранит значения одной колонки непрерывно и применяет к ним dictionary encoding и bit packing. JSON - строки смешанного содержимого, сжатие работает хуже. JSON-колонка с 1M строк занимает в 5–10 раз больше места, чем эквивалентные типизированные колонки Parquet.

Нет Predicate Pushdown. Когда Spark читает Parquet и встречает filter("amount > 1000") - он использует Min/Max статистику Row Groups и пропускает неподходящие блоки файла. JSON-строка внутри Parquet - непрозрачна: Spark не может применить pushdown к полям внутри JSON без предварительного парсинга.

Архитектурный компромисс: хранить сырой JSON в Bronze-слое (для воспроизводимости), парсить в типизированные структуры на Silver-слое (для производительности):

На каждом слое правило работы с JSON разное. Bronze: положить как есть, не парсить. Silver: распарсить в структуры, валидировать, обработать ошибки. Gold: работать только с типизированными колонками, JSON не должен появляться.


Чтение JSON-файлов vs JSON в строковых колонках

Важно различать две разные задачи:

1. Чтение JSON-файлов (.json, .jsonl, NDJSON) - файлы, где каждая строка это отдельный JSON-объект:

# Чтение JSON-файлов напрямую
df = spark.read.json("s3://bucket/events/*.jsonl")
# Spark автоматически определит схему и создаст StructType колонки

# С явной схемой (рекомендуется в production)
from pyspark.sql.types import StructType, StructField, StringType, LongType, DoubleType

schema = StructType([
    StructField("event_id",  StringType(),  True),
    StructField("user_id",   LongType(),    True),
    StructField("amount",    DoubleType(),  True),
    StructField("timestamp", StringType(),  True),
])
df = spark.read.schema(schema).json("s3://bucket/events/*.jsonl")

2. JSON как строковая колонка - Parquet/Kafka/Delta, где одна из колонок содержит JSON-строку:

# DataFrame где колонка 'payload' содержит JSON-строку
df = spark.read.parquet("s3://bucket/events/")
# schema: event_id STRING, received_at TIMESTAMP, payload STRING

df.printSchema()
# root
#  |-- event_id: string
#  |-- received_at: timestamp
#  |-- payload: string     ← это JSON-строка: '{"user_id":123,"amount":99.9}'

Второй случай - именно то, с чем мы работаем в большинстве реальных пайплайнов: Kafka-сообщения приходят как строки, сохраняются в Bronze как строки, и нужно их распарсить на Silver.


get_json_object: точечное извлечение через JSONPath

get_json_object(column, path) - извлекает одно поле из JSON-строки по JSONPath-выражению. Возвращает строку (всегда String, независимо от исходного типа).

from pyspark.sql import functions as F

df = spark.createDataFrame([
    ('{"user_id": 123, "amount": 99.9, "country": "RU"}',),
    ('{"user_id": 456, "amount": 250.0, "country": "DE"}',),
], ["payload"])

# Извлечение одного поля
df.select(
    F.get_json_object("payload", "$.user_id").alias("user_id"),
    F.get_json_object("payload", "$.amount").alias("amount"),
    F.get_json_object("payload", "$.country").alias("country"),
).show()
# user_id | amount | country
# 123     | 99.9   | RU
# 456     | 250.0  | DE

JSONPath синтаксис:

  • $.field - поле верхнего уровня
  • $.nested.field - вложенное поле
  • $.array[0] - первый элемент массива
  • $.array[*].field - поле во всех элементах массива
# Работа с вложенными структурами
payload = spark.createDataFrame([
    ('{"user": {"id": 1, "name": "Alice"}, "items": [{"sku": "A1", "qty": 2}]}',),
], ["payload"])

payload.select(
    F.get_json_object("payload", "$.user.id").alias("user_id"),
    F.get_json_object("payload", "$.user.name").alias("user_name"),
    F.get_json_object("payload", "$.items[0].sku").alias("first_sku"),
    F.get_json_object("payload", "$.items[0].qty").alias("first_qty"),
)

Главный недостаток get_json_object: каждый вызов парсирует всю JSON-строку заново. Пять вызовов get_json_object на одной строке - пять полных парсингов Jackson. На миллиардах строк это ощутимо.

Когда использовать: только если нужно извлечь 1–2 поля и остальное не интересует. Для большего числа полей - json_tuple или from_json.


json_tuple: несколько полей за один проход

json_tuple(column, field1, field2, ...) - Generator-функция, извлекающая несколько полей верхнего уровня за один парсинг JSON-строки:

from pyspark.sql import functions as F

df.select(
    F.json_tuple("payload", "user_id", "amount", "country")
    .alias("user_id", "amount", "country")
).show()

Ограничения json_tuple:

  • Только поля верхнего уровня - вложенные объекты не поддерживаются ($.user.id не работает)
  • Возвращает только String - нет автоматической конвертации типов
  • Нельзя использовать в Spark SQL напрямую (только через Python API как Generator)

Сравнение производительности для извлечения 5 полей из 100M строк:

Подход CPU time Примечание
5 × get_json_object ~180s 5 полных парсингов каждой строки
json_tuple (5 полей) ~45s 1 парсинг, 5 полей
from_json с StructType ~38s 1 парсинг + типизация

json_tuple в 4x быстрее get_json_object для множества полей. from_json ещё чуть быстрее за счёт оптимизаций в Catalyst.

Правило выбора:

  • 1 поле → get_json_object
  • 2–5 полей верхнего уровня, нужны строки → json_tuple
  • Больше 5 полей, нужна типизация, есть вложенность → from_json

from_json: полный парсинг в типизированную структуру

from_json(column, schema) - главный инструмент для серьёзной работы с JSON. Парсирует JSON-строку в структурированный объект Spark (StructType, ArrayType, вложенные структуры). Возвращает типизированные значения, а не строки.

from pyspark.sql import functions as F
from pyspark.sql.types import (
    StructType, StructField, StringType, LongType, DoubleType,
    ArrayType, TimestampType, BooleanType
)

# Определяем схему явно - рекомендуется в production
ORDER_SCHEMA = StructType([
    StructField("order_id",   StringType(),  False),
    StructField("user_id",    LongType(),    False),
    StructField("amount",     DoubleType(),  True),
    StructField("currency",   StringType(),  True),
    StructField("created_at", StringType(),  True),
    StructField("address", StructType([
        StructField("country", StringType(), True),
        StructField("city",    StringType(), True),
        StructField("zip",     StringType(), True),
    ]), True),
    StructField("items", ArrayType(StructType([
        StructField("sku",      StringType(), True),
        StructField("quantity", LongType(),   True),
        StructField("price",    DoubleType(), True),
    ])), True),
])

# Парсинг
parsed = df.select(
    F.col("event_id"),
    F.col("received_at"),
    F.from_json("payload", ORDER_SCHEMA).alias("order")
)

# Доступ к вложенным полям через dot notation
result = parsed.select(
    "event_id",
    "received_at",
    "order.order_id",
    "order.user_id",
    "order.amount",
    "order.address.country",
    "order.address.city",
)

После from_json колонка order имеет тип StructType - полноценную структуру Spark. Column Pruning работает: если из структуры нужны только order_id и amount, Spark читает только эти поля.

Доступ к вложенным полям

После from_json работают несколько способов доступа к полям:

# Dot notation (предпочтительно)
parsed.select("order.user_id", "order.address.country")

# Через col() с кавычками для имён с точками
parsed.select(F.col("order.user_id"), F.col("order.address.city"))

# Через getField (явный метод)
parsed.select(
    F.col("order").getField("user_id"),
    F.col("order").getField("address").getField("country")
)

# Раскрыть все поля структуры через .*
parsed.select("order.*")   # эквивалентно SELECT order.* FROM ...

Работа с массивами внутри JSON

Если JSON содержит массивы (список товаров в заказе, список событий), нужно explode() для нормализации:

from pyspark.sql import functions as F

# После from_json: колонка order.items - это ArrayType
# explode() превращает массив в отдельные строки
order_items = parsed.select(
    "order.order_id",
    "order.user_id",
    F.explode("order.items").alias("item")
).select(
    "order_id",
    "user_id",
    "item.sku",
    "item.quantity",
    "item.price",
    (F.col("item.quantity") * F.col("item.price")).alias("line_total")
)

order_items.show()
# order_id | user_id | sku  | quantity | price | line_total
# ORD-001  | 123     | A1   | 2        | 99.9  | 199.8
# ORD-001  | 123     | B3   | 1        | 49.9  | 49.9
# ORD-002  | 456     | C5   | 3        | 19.9  | 59.7

posexplode() добавит позицию элемента в массиве - полезно для сохранения порядка или соотнесения двух параллельных массивов.


schema_of_json: автоматический вывод схемы

Иногда схема JSON заранее неизвестна. schema_of_json(json_string) анализирует образец JSON и возвращает строку DDL со схемой:

from pyspark.sql import functions as F

sample = '{"user_id": 123, "amount": 99.9, "tags": ["vip", "active"], "address": {"city": "Moscow"}}'

# schema_of_json возвращает строку DDL
schema_ddl = spark.range(1) \
    .select(F.schema_of_json(F.lit(sample))) \
    .first()[0]

print(schema_ddl)
# STRUCT<address: STRUCT<city: STRING>, amount: DOUBLE, tags: ARRAY<STRING>, user_id: BIGINT>

# Использовать эту строку DDL в from_json
inferred_schema = schema_ddl
df.select(F.from_json("payload", inferred_schema).alias("data"))

Spark может также использовать schema_of_json с Python StructType:

# Более питоничный вариант через spark.read.json
sample_rdd = spark.sparkContext.parallelize([sample])
inferred_df = spark.read.json(sample_rdd)
print(inferred_df.schema)
# StructType с автоматически выведенными типами

Ограничения schema_of_json и когда это опасно

schema_of_json анализирует один образец. В production это создаёт серьёзные проблемы:

  • Необязательные поля: если в образце нет поля discount, схема его не включит. Строки с discount будут терять это поле при парсинге.
  • Неустойчивые типы: "amount": 100 → тип BIGINT. В реальных данных может быть "amount": 100.5 → тип DOUBLE. Несовпадение типа → NULL при парсинге.
  • Эволюция схемы: источник добавил новое поле - автовыведенная схема не знает о нём.
# ПЛОХО: вывод схемы из одной строки для production ETL
sample = df.limit(1).select("payload").first()[0]
schema = schema_of_json(sample)  # опасно!

# ХОРОШО: явная схема зафиксирована в коде и версионирована в git
ORDER_SCHEMA = StructType([...])  # явное определение

Допустимое использование schema_of_json: быстрое прототипирование в notebook, однократный анализ нового источника для последующей фиксации схемы, помощь в написании первоначальной схемы (скопировать DDL и превратить в StructType).


Обработка битых записей: режимы парсинга

from_json имеет три режима обработки ошибок, аналогично spark.read.json:

# Режим PERMISSIVE (по умолчанию): NULL для нераспарсенных полей
# Битые записи целиком → NULL структура
broken = spark.createDataFrame([
    ('{"order_id": "ORD-1", "amount": 99.9}',),   # корректный
    ('{broken json}',),                              # некорректный JSON
    ('{"order_id": "ORD-3", "amount": "not_a_number"}',),  # неверный тип
], ["payload"])

SCHEMA = StructType([
    StructField("order_id", StringType(), True),
    StructField("amount",   DoubleType(), True),
])

# PERMISSIVE: некорректный JSON → вся структура NULL
parsed = broken.select(F.from_json("payload", SCHEMA).alias("order"))
parsed.show()
# order
# {ORD-1, 99.9}
# {null, null}       ← битый JSON → вся структура NULL
# {ORD-3, null}      ← строка вместо числа → только amount = NULL

Опции for_json:

from pyspark.sql import functions as F

# PERMISSIVE + columnNameOfCorruptRecord: сохранить битую строку
SCHEMA_WITH_CORRUPT = StructType([
    StructField("order_id",        StringType(), True),
    StructField("amount",          DoubleType(), True),
    StructField("_corrupt_record", StringType(), True),  # обязательно в схеме!
])

parsed = broken.select(
    F.from_json("payload", SCHEMA_WITH_CORRUPT,
                options={"columnNameOfCorruptRecord": "_corrupt_record"}
    ).alias("order")
).select(
    "order.order_id",
    "order.amount",
    "order._corrupt_record"
)

# order_id | amount | _corrupt_record
# ORD-1    | 99.9   | null
# null     | null   | {broken json}    ← битая строка сохранена
# ORD-3    | null   | null             ← тип неверный, но JSON валиден

Разделение валидных и битых записей:

# Полный паттерн для production ETL
SCHEMA_WITH_CORRUPT = StructType([
    StructField("order_id",        StringType(), True),
    StructField("amount",          DoubleType(), True),
    StructField("_corrupt_record", StringType(), True),
])

all_parsed = kafka_df.select(
    F.col("offset"),
    F.col("timestamp"),
    F.from_json("value", SCHEMA_WITH_CORRUPT,
                {"columnNameOfCorruptRecord": "_corrupt_record"}
    ).alias("data")
).select("offset", "timestamp", "data.*")

# Разделить на валидные и битые
valid   = all_parsed.filter(F.col("_corrupt_record").isNull())
corrupt = all_parsed.filter(F.col("_corrupt_record").isNotNull())

# Записать в разные места
valid.drop("_corrupt_record") \
    .write.mode("append").parquet("s3://silver/orders/")

corrupt.write.mode("append").parquet("s3://quarantine/orders/")

to_json: генерация JSON из структур

to_json(column) - обратная операция: превращает StructType, ArrayType или MapType колонку в JSON-строку. Используется для:

  • Отправки данных в Kafka (Kafka ожидает строку)
  • Подготовки payload для внешних API
  • Сохранения сложных структур как JSON в Parquet (Anti-pattern, но иногда нужен)
from pyspark.sql import functions as F

# Простая структура → JSON
df = spark.createDataFrame([
    (1, "Alice", 99.9, ["tag1", "tag2"]),
    (2, "Bob",   150.0, ["vip"]),
], ["user_id", "name", "amount", "tags"])

# Создать StructType из нескольких колонок и сериализовать
df.select(
    F.to_json(
        F.struct("user_id", "name", "amount", "tags")
    ).alias("json_payload")
).show(truncate=False)
# {"user_id":1,"name":"Alice","amount":99.9,"tags":["tag1","tag2"]}
# {"user_id":2,"name":"Bob","amount":150.0,"tags":["vip"]}

Сериализация Map-структур:

# MapType → JSON object
df.select(
    F.to_json(
        F.create_map(
            F.lit("country"), F.col("country"),
            F.lit("city"),    F.col("city")
        )
    ).alias("location_json")
)
# {"country":"RU","city":"Moscow"}

Настройка форматирования:

# Опции to_json
df.select(
    F.to_json(
        F.struct("order_id", "amount", "created_at"),
        options={
            "timestampFormat": "yyyy-MM-dd'T'HH:mm:ss.SSSZ",
            "dateFormat": "yyyy-MM-dd",
            "ignoreNullFields": "true",    # не включать null-поля в JSON
        }
    ).alias("payload")
)

ignoreNullFields: true - полезная опция: если поле NULL, оно не попадёт в JSON. Это уменьшает размер payload и соответствует поведению большинства API (не включать поля с null).

Паттерн: Kafka producer из Spark:

# Подготовка Kafka-сообщений
kafka_messages = processed_df.select(
    F.col("user_id").cast("string").alias("key"),
    F.to_json(
        F.struct(
            "order_id", "user_id", "amount", "items", "created_at"
        ),
        options={"ignoreNullFields": "true"}
    ).alias("value")
)

kafka_messages.write \
    .format("kafka") \
    .option("kafka.bootstrap.servers", "broker:9092") \
    .option("topic", "order_events") \
    .save()

Эволюция схемы JSON: как справляться с изменениями

JSON-источники редко имеют стабильную схему. Поставщик API добавляет новые поля, меняет типы, убирает необязательные поля. Это называется Schema Evolution.

Стратегия 1: Широкая схема с nullable полями

# Зафиксировать максимально широкую схему
# Все поля nullable=True - новые поля игнорируются, отсутствующие → NULL
WIDE_SCHEMA = StructType([
    StructField("order_id",   StringType(),  True),
    StructField("amount",     DoubleType(),  True),
    StructField("currency",   StringType(),  True),   # было добавлено в v2
    StructField("discount",   DoubleType(),  True),   # было добавлено в v3
    StructField("created_at", StringType(),  True),
])

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

Стратегия 2: Хранить raw + парсить при чтении

# Bronze: хранить сырой JSON без изменений
kafka_df.select(
    F.col("offset"),
    F.col("timestamp").alias("ingested_at"),
    F.col("value").cast("string").alias("raw_json"),  # сырой JSON как строка
).write.mode("append").parquet("s3://bronze/events/")

# Silver: парсить при чтении с актуальной схемой
bronze = spark.read.parquet("s3://bronze/events/")

current_schema = load_schema_from_registry("order_events", version="latest")
silver = bronze.select(
    "offset",
    "ingested_at",
    F.from_json("raw_json", current_schema).alias("event")
)

Плюс: Bronze всегда воспроизводим - можно перепарсить с новой схемой. Минус: дополнительный pass по данным для парсинга.

Стратегия 3: schema_of_json + mergeSchema

# Автоматический вывод схемы из нескольких образцов
samples = df.limit(100).select("payload").collect()
json_strings = [row.payload for row in samples]

# Загрузить как JSON для вывода объединённой схемы
inferred_schema = spark.read.json(
    spark.sparkContext.parallelize(json_strings)
).schema

# Применить к основному DataFrame
df.select(F.from_json("payload", inferred_schema).alias("data"))

JSON и Parquet: физика хранения

Разберём подробно, что происходит когда JSON-данные попадают в Parquet.

Правильный способ: типизированные nested columns

# ХОРОШО: StructType в Parquet - полноценная вложенная структура
parsed_df = raw_df.select(
    F.from_json("payload", ORDER_SCHEMA).alias("order")
).select("order.*")   # развернуть в отдельные колонки

parsed_df.write.parquet("s3://silver/orders/")
# Parquet файл содержит:
# order_id: STRING column
# user_id: BIGINT column
# amount: DOUBLE column
# address.country: STRING column  ← nested column
# address.city: STRING column
# items.sku: STRING column        ← nested array column
# items.quantity: BIGINT column

В Parquet каждое поле StructType хранится как отдельная колонка со своей статистикой Min/Max и dictionary encoding. Column Pruning работает: если запросу нужен только amount, Spark прочитает только эту колонку.

Плохой способ: JSON-строка в Parquet

# ПЛОХО: сохранить сырой JSON как строку внутри Parquet
df.withColumn("json_data", F.col("payload")) \
  .write.parquet("s3://silver/events_bad/")
# Parquet файл содержит одну строковую колонку с JSON
# Никакой статистики по полям внутри JSON
# Нет Column Pruning для вложенных полей
# Нет Predicate Pushdown для условий на поля внутри JSON

Predicate Pushdown и JSON: что реально работает

Это ключевое ограничение: фильтрация по полям внутри JSON-строки в Parquet не использует Row Group статистику. Spark вынужден читать все Row Groups, парсить JSON в каждой строке, и только потом применять фильтр.

# Запрос к плохо организованному Silver
silver_bad = spark.read.parquet("s3://silver/events_with_json_string/")

# Filter по JSON-полю: Spark прочитает ВСЕ данные, потом отфильтрует
silver_bad \
    .filter(F.get_json_object("payload", "$.amount") > 1000) \
    .count()
# explain() покажет: нет PushedFilters, полный scan

# Запрос к правильно организованному Silver (StructType колонки)
silver_good = spark.read.parquet("s3://silver/events_structured/")
silver_good.filter(F.col("amount") > 1000).count()
# explain() покажет: PushedFilters: [IsNotNull(amount), GreaterThan(amount,1000.0)]

Column Pruning для вложенных структур: частично работает. Если структура раскрыта в отдельные Parquet-колонки, Spark читает только нужные. Если всё хранится как JSON-строка - читается вся строка всегда.


Полный ETL-пайплайн: Bronze → Silver

Соберём всё вместе в production-ready ETL:

from pyspark.sql import functions as F
from pyspark.sql.types import (
    StructType, StructField, StringType, LongType, DoubleType,
    ArrayType, TimestampType
)

# ===============================================================
# Схема фиксирована и версионирована в коде
# ===============================================================
ORDER_ITEM_SCHEMA = StructType([
    StructField("sku",      StringType(), True),
    StructField("quantity", LongType(),   True),
    StructField("price",    DoubleType(), True),
])

ADDRESS_SCHEMA = StructType([
    StructField("country", StringType(), True),
    StructField("city",    StringType(), True),
    StructField("zip",     StringType(), True),
])

ORDER_SCHEMA = StructType([
    StructField("order_id",        StringType(),            False),
    StructField("user_id",         LongType(),              False),
    StructField("amount",          DoubleType(),            True),
    StructField("currency",        StringType(),            True),
    StructField("status",          StringType(),            True),
    StructField("address",         ADDRESS_SCHEMA,          True),
    StructField("items",           ArrayType(ORDER_ITEM_SCHEMA), True),
    StructField("_corrupt_record", StringType(),            True),
])

# ===============================================================
# Чтение Bronze (сырые JSON-строки)
# ===============================================================
bronze = spark.read \
    .parquet("s3://bronze/orders/date=2026-05-16/") \
    .select("offset", "ingested_at", "raw_json")

# ===============================================================
# Парсинг с обработкой ошибок
# ===============================================================
parsed = bronze.select(
    "offset",
    "ingested_at",
    F.from_json(
        "raw_json",
        ORDER_SCHEMA,
        {"columnNameOfCorruptRecord": "_corrupt_record",
         "timestampFormat": "yyyy-MM-dd'T'HH:mm:ss.SSSZ"}
    ).alias("data")
)

# Разделить потоки
valid_df = parsed \
    .filter(F.col("data._corrupt_record").isNull()) \
    .select(
        "offset",
        "ingested_at",
        F.col("data.order_id").alias("order_id"),
        F.col("data.user_id").alias("user_id"),
        F.col("data.amount").alias("amount"),
        F.col("data.currency").alias("currency"),
        F.col("data.status").alias("status"),
        F.col("data.address.country").alias("country"),
        F.col("data.address.city").alias("city"),
        F.col("data.items").alias("items"),
    )

corrupt_df = parsed \
    .filter(F.col("data._corrupt_record").isNotNull()) \
    .select(
        "offset",
        "ingested_at",
        F.col("data._corrupt_record").alias("raw_json"),
        F.current_timestamp().alias("quarantined_at"),
    )

# ===============================================================
# Дополнительные трансформации
# ===============================================================
# Добавить метрики на уровне заказа
enriched = valid_df \
    .withColumn("item_count", F.size("items")) \
    .withColumn(
        "max_item_price",
        F.aggregate(
            "items",
            F.lit(0.0),
            lambda acc, item: F.greatest(acc, item.getField("price"))
        )
    ) \
    .withColumn(
        "created_date",
        F.to_date(F.col("ingested_at"))
    )

# ===============================================================
# Запись Silver (типизированные колонки в Parquet)
# ===============================================================
enriched.write \
    .mode("append") \
    .partitionBy("created_date", "country") \
    .parquet("s3://silver/orders/")

corrupt_df.write \
    .mode("append") \
    .parquet("s3://quarantine/orders/")

# ===============================================================
# Статистика прогона
# ===============================================================
total_raw  = bronze.count()
total_ok   = enriched.count()
total_bad  = corrupt_df.count()

print(f"Processed: {total_raw} | Valid: {total_ok} | Corrupt: {total_bad}")
print(f"Error rate: {total_bad / total_raw * 100:.2f}%")

Normalize nested arrays: пайплайн нормализации

Часто нужно из вложенного JSON создать нормализованную таблицу фактов. Пример: из заказа с массивом items создать строчку на каждый товар:

# После парсинга: каждый заказ - одна строка, items - массив
order_items = valid_df.select(
    "order_id",
    "user_id",
    "created_date",
    F.explode("items").alias("item")
).select(
    "order_id",
    "user_id",
    "created_date",
    F.col("item.sku").alias("sku"),
    F.col("item.quantity").alias("quantity"),
    F.col("item.price").alias("unit_price"),
    (F.col("item.quantity") * F.col("item.price")).alias("line_total"),
)

# Аналитика без UDF - только типизированные колонки
order_items.groupBy("sku") \
    .agg(
        F.sum("line_total").alias("total_revenue"),
        F.sum("quantity").alias("total_units"),
        F.avg("unit_price").alias("avg_price"),
        F.countDistinct("order_id").alias("order_count")
    ) \
    .orderBy(F.desc("total_revenue"))

Антипаттерны работы с JSON в Spark

1. Python UDF с json.loads() - самое медленное решение:

# ОЧЕНЬ ПЛОХО: Python UDF парсит JSON в Python runtime
import json
from pyspark.sql.types import DoubleType

@F.udf(returnType=DoubleType())
def extract_amount(payload):
    try:
        return json.loads(payload).get("amount")
    except:
        return None

df.withColumn("amount", extract_amount("payload"))
# Это в 5-10x медленнее from_json:
# - Данные сериализуются из JVM в Python
# - Python парсит JSON (медленнее Jackson)
# - Результат сериализуется обратно в JVM

2. Многократный get_json_object:

# ПЛОХО: каждый вызов парсит строку заново
df.select(
    F.get_json_object("payload", "$.a"),
    F.get_json_object("payload", "$.b"),
    F.get_json_object("payload", "$.c"),
    F.get_json_object("payload", "$.d"),
    F.get_json_object("payload", "$.e"),
)
# 5 парсингов каждой строки

# ХОРОШО: один from_json
df.select(F.from_json("payload", schema).alias("d")) \
  .select("d.a", "d.b", "d.c", "d.d", "d.e")

3. inferSchema на всех данных:

# ПЛОХО в production: inferSchema при чтении JSON-файлов
df = spark.read.json("s3://bucket/events/")   # inferSchema = двойной scan!

# ХОРОШО: явная схема
df = spark.read.schema(KNOWN_SCHEMA).json("s3://bucket/events/")

4. Аналитика по raw JSON-строке без парсинга:

# ПЛОХО: фильтр по JSON-строке - нет Pushdown, полный scan
df.filter(F.get_json_object("payload", "$.country") == "RU") \
  .agg(F.count("*"))

# ХОРОШО: распарсить в Silver, фильтровать по колонкам
silver = df.select(F.from_json("payload", schema).alias("d")).select("d.*")
silver.filter(F.col("country") == "RU").agg(F.count("*"))

5. Хранить JSON-строку в Parquet вместо структур:

Хранить payload STRING в Parquet вместо типизированных колонок - это убивает весь смысл колоночного хранения. Column Pruning и Predicate Pushdown не работают для полей внутри JSON-строки.


Чек-лист: работа с JSON в production

При проектировании схемы:

  • Зафиксировать схему явно в коде (StructType), не полагаться на inferSchema
  • Сделать все поля nullable=True для устойчивости к schema evolution
  • Добавить _corrupt_record поле для обработки битых записей
  • Включить в схему только поля, которые реально используются downstream

При парсинге:

  • Использовать from_json вместо множества get_json_object
  • Обрабатывать _corrupt_record - направлять в quarantine, не молча отбрасывать
  • Парсить один раз в начале pipeline, дальше работать только с типизированными колонками

При хранении:

  • Bronze: хранить raw JSON как строку (воспроизводимость)
  • Silver: StructType-колонки в Parquet (производительность, Column Pruning)
  • Gold: никакого JSON, только аналитические колонки

При чтении:

  • Проверить в explain(), что PushedFilters содержит условия фильтрации
  • Проверить ReadSchema - только нужные колонки читаются?
  • Не применять аналитические запросы к JSON-строкам - парсить сначала