Чтение и запись файлов в 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")