Чтение и запись JSON-файлов в PySpark

Однострочный и многострочный JSON, опции multiLine и samplingRatio, обработка повреждённых записей, сжатие и партиционирование при записи.

core

Введение

JSON-файлы — популярный выбор для структурированных данных. PySpark предлагает гибкие методы чтения и записи JSON-файлов с различными опциями конфигурации для разных задач.

Шаг 1: Инициализация SparkSession

from pyspark.sql import SparkSession

# Инициализация SparkSession
spark = SparkSession.builder \
    .appName("JSON Handling in PySpark") \
    .getOrCreate()

print("Spark is ready to handle JSON files! 🚀")

Шаг 2: Чтение JSON-файла 📥

Базовое чтение JSON

Для стандартных JSON-файлов, где каждая строка содержит отдельный JSON-объект:

# Чтение по умолчанию: каждая строка должна быть отдельным JSON-объектом
# { "id": 1, "name": "Alice" }
# { "id": 2, "name": "Bob" }
df = spark.read.json("/path/to/your/data.json")
df.show(5)

Многострочные JSON-файлы

Для форматированных JSON-массивов или объектов, занимающих несколько строк:

# multiline.json: отформатированный JSON-массив
# [
#   { "id": 1,
#     "name": "Alice",
#     "dept": "Engineering"
#   },
#   { "id": 2,
#     "name": "Bob",
#     "dept": "Finance"
#   }
# ]

df = spark.read \
    .option("multiLine", True) \
    .json("/path/to/multiline.json")

df.show()

Основные опции чтения

  • Multiline: установите True для записей, занимающих несколько строк
  • InferSchema: автоматически определяет типы данных из содержимого
  • SamplingRatio: управляет долей строк, выборочно проверяемых для вывода схемы (от 0 до 1)
# Чтение JSON с многострочным режимом и выборкой схемы
df = spark.read.option("multiline", True) \
               .option("samplingRatio", 0.5) \
               .json("/path/to/your/multiline_data.json")

df.show(5)

Обработка некорректных данных

Опция mode определяет поведение при обнаружении повреждённых записей:

  • PERMISSIVE (по умолчанию): помещает их в колонку _corrupt_record
  • DROPMALFORMED: удаляет некорректные записи
  • FAILFAST: останавливает чтение и выбрасывает ошибку при первой повреждённой записи
# Обработка повреждённых JSON-записей
df = spark.read.option("mode", "FAILFAST") \
               .json("/path/to/your/data.json")

df.show(5)

Шаг 3: Запись DataFrame в JSON-файл 📤

Базовая запись JSON

# Простая запись в JSON
df.write.json("/path/to/output_data")

Основные опции записи

Mode управляет поведением при наличии существующих файлов:

  • overwrite: заменяет существующие файлы
  • append: добавляет к существующим данным
  • ignore: пропускает запись, если файлы существуют
  • error / errorifexists: выбрасывает исключение

Compression (сжатие): gzip, bzip2, lz4, snappy, deflate.

# Запись JSON с режимом overwrite и сжатием gzip
df.write.mode("overwrite") \
        .option("compression", "gzip") \
        .json("/path/to/compressed_output")

Партиционирование выходных данных

# Запись JSON с партиционированием по году
df.write.option("header", True) \
        .partitionBy("year") \
        .json("/path/to/partitioned_output")

Итог

PySpark упрощает работу с JSON-файлами благодаря гибким опциям вывода схемы, обработки многострочного формата, сжатия и партиционирования данных — необходимым возможностям для эффективного управления большими структурированными наборами данных.