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

Универсальные паттерны чтения и записи данных: CSV, JSON, Parquet, Hive-таблицы, JDBC. Режимы чтения (PERMISSIVE, DROPMALFORMED, FAILFAST) и режимы записи (overwrite, append, ignore, error).

core

Чтение файлов в PySpark

Базовый паттерн чтения файлов в PySpark:

from pyspark.sql import SparkSession

# Шаг 1: Инициализация SparkSession
spark = SparkSession.builder \
    .appName("Read File Example") \
    .getOrCreate()

# Шаг 2: Чтение данных
df = spark.read \
    .format("<file_format>") \
    .option("<option_name>", "<option_value>") \
    .load("<path_to_file_or_folder>")

# Шаг 3: Просмотр данных
df.show(5)
df.printSchema()

Параметры чтения

Параметр Описание Пример
<file_format> Тип источника данных "csv", "json", "parquet", "orc", "jdbc"
option() Опции чтения .option("header", "true")
load() Путь к файлу или директории "/datasets/customers.csv"

Чтение CSV-файлов

df = spark.read \
    .format("csv") \
    .option("header", "true") \
    .option("inferSchema", "true") \
    .load("/datasets/customers.csv")

Чтение Parquet-файлов

df = spark.read \
    .format("parquet") \
    .load("/datasets/sales.parquet")

Чтение JSON-файлов

df = spark.read \
    .format("json") \
    .load("/datasets/events.json")

Чтение из Hive-таблиц

# Чтение Hive-таблицы customers из схемы default в DataFrame
df_hive = spark.read.table("default.customers")

df_hive.show(5)
df_hive.printSchema()

Через Spark SQL:

# Запрос к Hive-таблице через SQL
df_hive_sql = spark.sql("SELECT id, name, email FROM default.customers WHERE country='US'")
df_hive_sql.show(5)

Чтение из базы данных (JDBC)

df = spark.read \
    .format("jdbc") \
    .option("url", "jdbc:postgresql://localhost:5432/mydb") \
    .option("dbtable", "public.customers") \
    .option("user", "postgres") \
    .option("password", "mypassword") \
    .load()

Краткие формы для популярных форматов

df_csv = spark.read.csv(path, header=True, inferSchema=True)
df_json = spark.read.json(path)
df_parquet = spark.read.parquet(path)

Режимы чтения

Режим Описание
PERMISSIVE Режим по умолчанию: поля повреждённых записей устанавливаются в null, сырые данные сохраняются в колонке _corrupt_record
DROPMALFORMED Удаляет целые строки при любом несоответствии схеме или повреждении, загружая только корректные записи
FAILFAST Немедленно выбрасывает исключение при обнаружении первой некорректной записи

Задаётся через .option("mode", "PERMISSIVE") или .mode("DROPMALFORMED").


Запись файлов в PySpark

Общий паттерн записи файлов в PySpark:

from pyspark.sql import SparkSession

# Шаг 1: Инициализация SparkSession
spark = SparkSession.builder \
    .appName("Write File Example") \
    .getOrCreate()

# Шаг 2: Запись данных
df.write \
    .format("<file_format>") \
    .option("<option_name>", "<option_value>") \
    .mode("<save_mode>") \
    .save("<path_to_output_file_or_folder>")

Параметры записи

Параметр Описание Пример
<file_format> Тип выходных данных "csv", "json", "parquet", "orc", "jdbc"
option() Опции записи .option("header", "true")
mode() Поведение при наличии существующих данных "overwrite", "append", "ignore", "error"
save() Путь к выходному файлу или директории "/output/customers_parquet"

Режимы сохранения

Режим Описание
overwrite Заменяет существующие данные
append Добавляет данные к существующим файлам или таблицам
ignore Пропускает запись, если место назначения уже существует
error / errorifexists Завершается ошибкой, если место назначения существует

Запись CSV-файлов

df.write \
    .format("csv") \
    .option("header", "true") \
    .mode("overwrite") \
    .save("/output/customers_csv")

Запись Parquet-файлов

df.write \
    .format("parquet") \
    .mode("overwrite") \
    .save("/output/sales_parquet")

Запись JSON-файлов

df.write \
    .format("json") \
    .mode("overwrite") \
    .save("/output/events_json")

Запись в Hive-таблицы

# Запись DataFrame в Hive-таблицу в схеме default
df.write \
    .mode("overwrite") \
    .saveAsTable("default.customers_copy")

Запись в базу данных (JDBC)

df.write \
    .format("jdbc") \
    .option("url", "jdbc:postgresql://localhost:5432/mydb") \
    .option("dbtable", "public.customers") \
    .option("user", "postgres") \
    .option("password", "mypassword") \
    .mode("append") \
    .save()

Краткие формы для популярных форматов

df.write.csv(path="/output/customers", header=True, mode="overwrite")
df.write.json(path="/output/events", mode="overwrite")
df.write.parquet(path="/output/sales", mode="overwrite")