Чтение и запись 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 позволяют быстро адаптировать настройки для масштабного чтения и записи данных — будь то заголовки, разделители, автоматическое определение схемы или сжатие данных.