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

Опции чтения CSV: header, inferSchema, sep, nullValue, dateFormat, режимы обработки ошибок. Опции записи: header, mode, delimiter, compression.

core

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

from pyspark.sql import SparkSession

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

print("Spark is ready! 🚀")

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

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

# Простое чтение CSV
df = spark.read.csv("/path/to/your/data.csv")
df.show(5)

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

  • Header: если CSV содержит строку заголовка, установите header=True, чтобы использовать её как имена колонок.
  • InferSchema: для автоматического определения типов данных каждой колонки установите inferSchema=True.
  • Delimiter: для использования нестандартного разделителя (например, ; или |) задайте опцию sep.
# Чтение CSV с опциями
df = spark.read.option("header", True) \
               .option("inferSchema", True) \
               .option("sep", ";") \
               .csv("/path/to/your/data.csv")

df.show(5)

Дополнительные опции чтения

  • Null-значения: используйте nullValue для указания строки, которая интерпретируется как null, например "N/A" или "NULL".
  • Формат даты: задайте dateFormat для указания формата дат в полях типа дата.
  • Режим: определяет, как Spark обрабатывает повреждённые записи:
  • PERMISSIVE (по умолчанию): помещает повреждённые записи в отдельную колонку.
  • DROPMALFORMED: удаляет строки с некорректными записями.
  • FAILFAST: выбрасывает ошибку при обнаружении повреждённой записи.
# Чтение CSV с расширенными опциями
df = spark.read.option("header", True) \
               .option("inferSchema", True) \
               .option("nullValue", "N/A") \
               .option("dateFormat", "MM/dd/yyyy") \
               .option("mode", "DROPMALFORMED") \
               .csv("/path/to/your/data.csv")

df.show(5)

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

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

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

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

  • Header: установите header=True, чтобы включить имена колонок в первую строку.
  • Mode: управляет поведением при наличии существующего файла или директории:
  • overwrite: заменяет существующие данные.
  • append: добавляет данные к существующему файлу.
  • ignore: пропускает запись, если файл уже существует.
  • error / errorifexists: выбрасывает ошибку, если файл существует.
# Запись CSV с опциями
df.write.option("header", True) \
        .mode("overwrite") \
        .csv("/path/to/output_data")

Дополнительные опции записи

  • Delimiter: задайте sep для указания нестандартного разделителя.
  • Compression: выберите метод сжатия: gzip, bzip2, lz4, snappy или deflate.
# Запись CSV с разделителем и сжатием
df.write.option("header", True) \
        .option("sep", ";") \
        .option("compression", "gzip") \
        .csv("/path/to/compressed_data")

Итог

Опции обработки CSV в PySpark позволяют быстро адаптировать настройки для масштабного чтения и записи данных — будь то заголовки, разделители, автоматическое определение схемы или сжатие данных.