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.
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-строкам - парсить сначала