Блок 3 - DataFrame API и SQL
20 вопросов: чтение форматов, join-типы, window functions, null-обработка, complex types и actions.
Как прочитать данные из JSON/CSV/Parquet?¶
# Parquet - колоночный формат, поддерживает predicate pushdown
df = spark.read.parquet("s3://bucket/data/")
# CSV - задаём схему вручную для надёжности
from pyspark.sql.types import StructType, IntegerType, StringType
schema = StructType().add("id", IntegerType()).add("name", StringType())
df = spark.read.option("header", "true").schema(schema).csv("data.csv")
# JSON Lines (один объект на строку)
df = spark.read.json("data.json")
# Многострочный JSON
df = spark.read.option("multiLine", "true").json("data.json")
Режимы обработки ошибок: PERMISSIVE (NULL для плохих полей, default), DROPMALFORMED (удалить строку), FAILFAST (остановить Job).
Типы Join - inner, left, outer, anti, semi¶
# inner - только совпадающие строки
df = a.join(b, "id", "inner")
# left (left_outer) - все строки A, NULL для B при несовпадении
df = a.join(b, "id", "left")
# full (outer) - все строки A и B
df = a.join(b, "id", "outer")
# left_semi - строки A, для которых есть совпадение в B (не берёт колонки B)
df = a.join(b, "id", "left_semi")
# left_anti - строки A, для которых НЕТ совпадения в B (инкрементальная загрузка)
new_only = new.join(existing, "id", "left_anti")
| Тип | Строки из A | Строки из B |
|---|---|---|
| inner | только совпадения | только совпадения |
| left | все | только совпадения |
| full | все | все |
| semi | совпадения | - (не включаются) |
| anti | несовпадения | - |
Какие физические стратегии Join использует Spark (SMJ, BHJ, SHJ)?¶
Это физические алгоритмы выполнения join, которые Catalyst выбирает в зависимости от размера таблиц и конфигурации:
| Стратегия | Полное название | Как работает | Когда выбирается |
|---|---|---|---|
| BHJ | Broadcast Hash Join | Маленькая таблица копируется на все Executor'ы, shuffle не нужен | Одна сторона < autoBroadcastJoinThreshold (10 МБ) |
| SMJ | Sort-Merge Join | Обе стороны shuffleются по ключу, сортируются, сливаются merge-итератором | По умолчанию для больших таблиц |
| SHJ | Shuffle Hash Join | Shuffle по ключу, build hash table для меньшей стороны, probe с большой | Одна сторона значительно меньше другой, не влезает в broadcast |
# Явные hints для управления стратегией
from pyspark.sql.functions import broadcast
big.join(broadcast(small), "id") # → BHJ (без shuffle)
big.hint("shuffle_hash").join(small, "id") # → SHJ
big.hint("merge").join(small, "id") # → SMJ
# Проверить в explain()
df.explain()
# == Physical Plan ==
# *(5) BroadcastHashJoin [id#1], [id#2], Inner, BuildRight
# или
# *(5) SortMergeJoin [id#1], [id#2], Inner
BHJ - самый быстрый (нет shuffle). SMJ - самый надёжный (любой размер, spill на диск). SHJ - промежуточный: быстрее SMJ, но требует больше памяти чем BHJ.
В чём разница между where() и filter()?¶
Разницы нет - это алиасы одного метода. df.filter(...) и df.where(...) полностью идентичны. Оба принимают Column expression или строку SQL:
df.filter("age > 18")
df.filter(col("age") > 18)
df.where("age > 18") # То же самое
Catalyst оптимизирует оба одинаково - применяет predicate pushdown где возможно.
Как создать колонку на основе условия (when, otherwise)?¶
from pyspark.sql.functions import when, col
# Простое условие
df = df.withColumn("category",
when(col("age") < 18, "minor")
.when(col("age") < 65, "adult")
.otherwise("senior"))
# Вложенные условия
df = df.withColumn("risk",
when((col("amount") > 1000) & (col("country") == "RU"), "high")
.when(col("amount") > 500, "medium")
.otherwise("low"))
when не имеет ограничений по вложенности. otherwise обязателен если нужен fallback - без него незаполненные строки получат null.
Что такое Window Functions? Примеры row_number, rank¶
Window Functions вычисляют значения в рамках группы строк (окна) без уменьшения числа строк (в отличие от groupBy).
from pyspark.sql.window import Window
from pyspark.sql.functions import row_number, rank, dense_rank, lag, lead, sum as spark_sum
w = Window.partitionBy("dept").orderBy("salary")
df = df.withColumn("rn", row_number().over(w)) # Уникальный, без пропусков
df = df.withColumn("rank", rank().over(w)) # С пропусками при одинаковых
df = df.withColumn("dense_rank", dense_rank().over(w)) # Без пропусков
df = df.withColumn("prev_salary", lag("salary", 1).over(w))
df = df.withColumn("next_salary", lead("salary", 1).over(w))
# Нарастающий итог (Running Total) - window frame
w_cum = Window.partitionBy("dept").orderBy("date").rowsBetween(Window.unboundedPreceding, 0)
df = df.withColumn("cumulative_sum", spark_sum("amount").over(w_cum))
Как работают groupBy и agg?¶
groupBy группирует строки по ключу. agg применяет агрегатные функции к каждой группе.
from pyspark.sql.functions import count, sum, avg, max, min, countDistinct, collect_list
df.groupBy("dept").agg(
count("*").alias("total_rows"),
sum("salary").alias("total_salary"),
avg("salary").alias("avg_salary"),
max("salary").alias("max_salary"),
countDistinct("employee_id").alias("unique_employees"),
collect_list("name").alias("names_list")
)
rollup и cube - расширения для иерархических и многомерных агрегаций (аналоги ROLLUP/CUBE в SQL).
Как удалить дубликаты: dropDuplicates vs distinct?¶
distinct() - удаляет строки, которые полностью идентичны по всем колонкам. dropDuplicates(subset) - удаляет дубли по выбранным колонкам, оставляя произвольную строку из группы.
df.distinct() # Все колонки
df.dropDuplicates(["user_id"]) # По ключу, оставить одну строку
df.dropDuplicates(["user_id", "date"]) # По составному ключу
Для выбора конкретной строки (самой свежей, с максимальным значением) используют Window + row_number:
w = Window.partitionBy("user_id").orderBy(col("updated_at").desc())
df.withColumn("rn", row_number().over(w)).filter("rn = 1").drop("rn")
Как обрабатывать пропущенные значения (Null)?¶
# Удалить строки с null в любой колонке
df.dropna()
df.dropna(how="all") # Только если все колонки null
df.dropna(subset=["id", "name"]) # Null в конкретных колонках
# Заполнить null значениями
df.fillna(0) # Все числовые колонки
df.fillna({"age": 0, "city": "N/A"}) # Разные значения по колонкам
df.fillna(df.agg({"salary": "mean"}).collect()[0][0], subset=["salary"]) # Средним
# Проверить наличие null
from pyspark.sql.functions import col, when, count
df.select([count(when(col(c).isNull(), c)).alias(c) for c in df.columns])
coalesce(col_a, col_b, default) - возвращает первое не-null значение из списка.
Как выполнить сырой SQL-запрос в PySpark?¶
# Зарегистрировать DataFrame как временную вью
df.createOrReplaceTempView("users")
# Выполнить SQL
result = spark.sql("""
SELECT dept, AVG(salary) as avg_salary
FROM users
WHERE active = true
GROUP BY dept
HAVING AVG(salary) > 50000
ORDER BY avg_salary DESC
""")
# Global temp view (видна в других сессиях в том же приложении)
df.createOrReplaceGlobalTempView("users_global")
spark.sql("SELECT * FROM global_temp.users_global")
Разница между select() и selectExpr()¶
select() принимает Column objects или строки-имена колонок. selectExpr() принимает SQL-выражения в виде строк - позволяет использовать SQL синтаксис в Python-коде.
# select - Column expressions
df.select(col("name"), (col("salary") * 1.1).alias("new_salary"))
# selectExpr - SQL strings
df.selectExpr("name", "salary * 1.1 AS new_salary", "UPPER(city) AS city_upper")
selectExpr удобен когда выражение проще написать на SQL чем через Column API.
Как переименовать колонку?¶
df.withColumnRenamed("old_name", "new_name")
# Несколько колонок
df.withColumnsRenamed({"old_a": "new_a", "old_b": "new_b"}) # Spark 3.4+
# Через select + alias
df.select(col("old_name").alias("new_name"), "*")
# Через toDF() - переименовать все колонки сразу
df.toDF("col1", "col2", "col3")
Как изменить тип данных (cast)?¶
from pyspark.sql.functions import col
from pyspark.sql.types import IntegerType, DoubleType, TimestampType
df = df.withColumn("age", col("age").cast(IntegerType()))
df = df.withColumn("amount", col("amount").cast(DoubleType()))
df = df.withColumn("ts", col("ts_string").cast(TimestampType()))
# Shorthand string syntax
df = df.withColumn("price", col("price").cast("double"))
df = df.withColumn("qty", col("qty").cast("int"))
При невозможности конвертации (например, буква → int) Spark вернёт null (в режиме PERMISSIVE) или бросит исключение (FAILFAST).
Как работать со сложными структурами (Array, Map, Struct)?¶
from pyspark.sql.functions import (
array, array_contains, size, explode, flatten,
map_keys, map_values, col, struct, get_json_object
)
# Array - доступ по индексу, contains, size
df.select(col("tags")[0], size(col("tags")), array_contains(col("tags"), "python"))
# Map - ключи и значения
df.select(map_keys("props"), map_values("props"), col("props")["color"])
# Struct - доступ через точку
df.select("address.city", "address.zip")
# Создать struct
df.select(struct("name", "age").alias("person"))
# Higher-order functions (Spark 3.x)
from pyspark.sql.functions import transform, filter as array_filter
df.withColumn("doubled", transform("scores", lambda x: x * 2))
df.withColumn("high_scores", array_filter("scores", lambda x: x > 90))
Что делают explode() и flatten()?¶
explode(col) - разворачивает массив или map: каждый элемент становится отдельной строкой. posexplode добавляет индекс. flatten - сплющивает массив массивов.
# explode: одна строка с array [1,2,3] → три строки
df.withColumn("item", explode(col("items_array")))
# posexplode: добавляет порядковый номер
df.select("id", posexplode("items").alias("pos", "item"))
# explode_outer: сохраняет строки с пустым/null массивом (→ null)
df.withColumn("item", explode_outer(col("items_array")))
# flatten: [[1,2],[3,4]] → [1,2,3,4]
df.withColumn("flat", flatten(col("nested_array")))
Как отфильтровать данные по регулярному выражению?¶
from pyspark.sql.functions import col, regexp_extract, regexp_replace
# Фильтр по regex - строки, где email заканчивается на .ru
df.filter(col("email").rlike(r"\.ru$"))
# Извлечь подстроку по regex
df.withColumn("domain", regexp_extract(col("email"), r"@(.+)$", 1))
# Заменить по regex
df.withColumn("clean", regexp_replace(col("text"), r"[^a-zA-Z0-9]", "_"))
Как получить схему DataFrame?¶
df.printSchema() # Вывод в консоль (дерево)
df.schema # StructType объект
df.dtypes # Список [("col_name", "col_type"), ...]
# Создать DF из схемы
from pyspark.sql.types import StructType, StructField, StringType, IntegerType
schema = StructType([
StructField("id", IntegerType(), nullable=False),
StructField("name", StringType(), nullable=True),
])
df = spark.read.schema(schema).csv("data.csv")
Явная схема важнее inferSchema=true: надёжнее, быстрее (не нужен дополнительный проход по данным).
Разница между collect(), take(), show() и first()¶
| Метод | Возвращает | Объём | Риск |
|---|---|---|---|
collect() |
List[Row] (всё) | Все данные | OOM на Driver |
take(N) |
List[Row] (первые N) | N строк | Безопасен |
first() |
Row (первая) | 1 строка | Безопасен |
show(N) |
None (печать) | N строк | Безопасен |
head(N) |
List[Row] | N строк | Аналог take |
В production collect() допустим только если данные заведомо малы (результат агрегации, lookup-таблица).
Когда использовать toPandas() и какие риски?¶
toPandas() - Action, который собирает все данные из DataFrame на Driver и конвертирует в pandas DataFrame. Риски: тот же, что у collect() - OOM если данных много. Используйте только:
- После агрегации, когда результат заведомо мал
- Для визуализации/отчётов на небольших выборках
- В конце pipeline, для передачи данных в ML-библиотеку
# Плохо: огромный DataFrame → OOM
pd_df = spark_df.toPandas()
# Хорошо: сначала агрегация или limit
pd_df = spark_df.groupBy("dept").agg(avg("salary")).toPandas()
pd_df = spark_df.limit(10000).toPandas()
Как добавить текущую дату/время в DataFrame?¶
from pyspark.sql.functions import current_date, current_timestamp, lit
df = df.withColumn("load_date", current_date()) # DATE
df = df.withColumn("load_ts", current_timestamp()) # TIMESTAMP
# Фиксированная дата
df = df.withColumn("snapshot_date", lit("2024-01-01").cast("date"))
Как объединить DataFrame вертикально: union vs unionByName¶
union() - объединяет по позиции колонок (как SQL UNION ALL). Схемы должны совпадать по порядку и типам, имена колонок игнорируются. unionByName() - объединяет по именам колонок, допускает разный порядок.
# union - совпадение по позиции
df_all = df1.union(df2)
# unionByName - совпадение по имени (безопаснее)
df_all = df1.unionByName(df2)
# Spark 3.x: allowMissingColumns=True - заполняет null для отсутствующих колонок
df_all = df1.unionByName(df2, allowMissingColumns=True)
Как Spark обрабатывает Null-значения при выполнении Join?¶
В SQL и Spark NULL != NULL - NULL означает "неизвестное значение", оно не равно ничему, включая другой NULL. Это напрямую влияет на результат join:
# Inner join: строки с NULL в join-ключе ОТБРАСЫВАЮТСЯ
a.join(b, "id") # Строки где id IS NULL из обеих таблиц - не попадают в результат
# Left join: строки с NULL id из A сохраняются, b-поля → NULL
a.join(b, "id", "left")
# Явная обработка NULL как равных (нестандартное поведение)
a.join(b, a.id.eqNullSafe(b.id)) # NULL == NULL → True
a.join(b, (a.id == b.id) | (a.id.isNull() & b.id.isNull()))
-- В SQL: по умолчанию NULL != NULL
SELECT * FROM a JOIN b ON a.id = b.id -- NULL-строки не войдут
-- Явно обработать NULL
SELECT * FROM a JOIN b ON a.id <=> b.id -- NULL-safe equals
Практическое правило: если join-ключ может содержать NULL - проверьте это заранее через df.filter(col("id").isNull()).count() и решите, нужно ли сохранять такие строки.
Почему Sort-Merge Join является стратегией по умолчанию для больших таблиц?¶
SMJ - безопасный выбор по умолчанию для больших данных: не требует загружать одну из сторон целиком в память.
Алгоритм: (1) Shuffle обеих таблиц по join-ключу → строки с одинаковым ключом на одном Executor'е; (2) Sort каждой стороны; (3) Merge-iterator: линейный проход по двум отсортированным потокам.
Потребление памяти: O(1) для merge-шага, spill на диск при необходимости - гарантии от OOM.
Broadcast Hash Join (BHJ) требует, чтобы меньшая сторона целиком поместилась в память каждого Executor'а. При ошибке (превышении порога) - риск OOM или тихий откат на SMJ без предупреждения.
Catalyst выбирает BHJ если table_size < spark.sql.autoBroadcastJoinThreshold (10 МБ по умолчанию). Для всего остального - SMJ как надёжный fallback.
В каких случаях Cross Join (Cartesian Product) оправдан?¶
Cross Join создаёт декартово произведение: каждая строка из A × каждая строка из B. Результат: N × M строк.
Без явного crossJoin() Spark бросает AnalysisException если нет join-условия - защита от случайного Cartesian.
Легитимные случаи:
# 1. Генерация всех комбинаций (рекомендательные системы)
categories = spark.createDataFrame([("electronics",), ("books",)], ["cat"])
users_x_cats = users.crossJoin(categories) # каждый пользователь × каждая категория
# 2. Broadcast маленькой таблицы-констант на весь DataFrame
thresholds = spark.createDataFrame([(100, "low"), (500, "high")], ["threshold", "tier"])
result = data.crossJoin(broadcast(thresholds)) # 2-строчная таблица
# 3. SQL: явный CROSS JOIN
spark.sql("SELECT * FROM users CROSS JOIN date_spine WHERE date BETWEEN start AND end")
Опасность: при N=1M и M=1M результат - 1 триллион строк. Всегда оценивайте размер до выполнения.
Как реализовать SCD Type 2 средствами Spark (Delta Lake)?¶
SCD Type 2 хранит историю изменений: при каждом обновлении старая запись "закрывается" (устанавливается effective_to), вставляется новая актуальная.
from delta.tables import DeltaTable
from pyspark.sql.functions import current_timestamp, lit
existing = DeltaTable.forPath(spark, "s3://bucket/dim_users/")
updates = spark.table("source_updates")
# 1. Добавить технические поля к обновлениям
updates = updates \
.withColumn("effective_from", current_timestamp()) \
.withColumn("effective_to", lit(None).cast("timestamp")) \
.withColumn("is_current", lit(True))
# 2. MERGE: закрыть старую запись + вставить новую
existing.alias("tgt").merge(
updates.alias("src"),
"tgt.user_id = src.user_id AND tgt.is_current = true"
).whenMatchedUpdate(set={
"effective_to": "src.effective_from",
"is_current": "false"
}).whenNotMatchedInsert(values={
"user_id": "src.user_id",
"name": "src.name",
"effective_from": "src.effective_from",
"effective_to": "null",
"is_current": "true"
}).execute()
Аналог в Iceberg: spark.sql("MERGE INTO target USING source ON target.user_id = source.user_id AND target.is_current = true ...")
Разница между внутренними (managed) и внешними (external) таблицами в Spark SQL?¶
| Managed (внутренняя) | External (внешняя) | |
|---|---|---|
| Управление данными | Spark владеет файлами | Пользователь управляет |
DROP TABLE |
Удаляет метаданные и данные | Удаляет только метаданные |
| Расположение данных | spark.sql.warehouse.dir |
Указанный LOCATION |
| Использование | Временные/промежуточные таблицы | Общие данные, несколько движков |
-- Внешняя таблица (данные сохраняются при DROP)
CREATE EXTERNAL TABLE events (id INT, event STRING)
STORED AS PARQUET
LOCATION 's3://bucket/data/events/';
-- Внутренняя (данные удалятся вместе с таблицей)
CREATE TABLE summary AS
SELECT dept, COUNT(*) AS cnt FROM employees GROUP BY dept;
Для production data lake рекомендуются external-таблицы - данные независимы от жизненного цикла metastore.
Как работает Partitioning на уровне файловой системы?¶
Партиционирование создаёт иерархию директорий в файловой системе по значениям выбранных колонок. Spark при чтении читает только нужные директории (partition pruning), не открывая лишние файлы.
# Запись с партиционированием по дате и стране
df.write \
.partitionBy("date", "country") \
.parquet("s3://bucket/events/")
# На диске:
# events/date=2024-01-01/country=US/part-0000.parquet
# events/date=2024-01-01/country=RU/part-0001.parquet
# events/date=2024-01-02/country=US/part-0002.parquet
# При чтении с фильтром - Spark читает только нужную директорию
df = spark.read.parquet("s3://bucket/events/") \
.filter("date = '2024-01-01' AND country = 'US'")
# Читается: events/date=2024-01-01/country=US/ - только 1 директория из 3
Правила хорошего партиционирования:
- Ключ партиции - низкая кардинальность: дата, страна, статус (не user_id)
- Размер одной партиции - 100 МБ–1 ГБ. Слишком мелкие → проблема Small Files
- Не более 3–4 уровней вложенности - иначе overhead на LIST операции
Какие существуют Write Modes (Append, Overwrite, ErrorIfExists, Ignore)?¶
Write Mode определяет поведение Spark при записи в существующий путь/таблицу:
| Mode | Поведение | Применение |
|---|---|---|
| append | Добавить новые файлы к существующим | Инкрементальная загрузка, streaming sink |
| overwrite | Удалить существующие данные, записать новые | Перезапись партиции, batch reload |
| errorIfExists | Бросить ошибку если путь уже существует | Default - защита от случайной перезаписи |
| ignore | Ничего не делать если путь существует | Idempotent "записать если пусто" |
# append - добавить к существующим данным
df.write.mode("append").parquet("s3://bucket/output/")
# overwrite - полностью заменить
df.write.mode("overwrite").parquet("s3://bucket/output/")
# Частичная перезапись партиции (только указанная партиция)
spark.conf.set("spark.sql.sources.partitionOverwriteMode", "dynamic")
df.write.mode("overwrite").partitionBy("date").parquet("s3://bucket/output/")
# Перезапишет только партиции, присутствующие в df - остальные не тронет
# ignore - если путь существует, пропустить запись
df.write.mode("ignore").parquet("s3://bucket/output/")
Dynamic Partition Overwrite (partitionOverwriteMode = dynamic) - важный паттерн: overwrite только конкретных партиций, а не всей таблицы.
Как реализовать идемпотентность при записи данных в Spark?¶
Идемпотентная запись - повторный запуск pipeline с теми же входными данными даёт тот же результат без дублирования.
Основные техники:
1. Partition overwrite - перезаписать только партицию текущего периода:
spark.conf.set("spark.sql.sources.partitionOverwriteMode", "dynamic")
df.write.mode("overwrite").partitionBy("date").parquet("s3://bucket/output/")
# При повторном запуске: та же партиция date=2024-01-01 перезапишется - дублей нет
2. MERGE INTO - upsert в Delta/Iceberg:
# При повторном запуске: update существующих + insert новых (не duplicate)
delta_table.alias("t").merge(new_data.alias("s"), "t.id = s.id") \
.whenMatchedUpdate(set={"value": "s.value"}) \
.whenNotMatchedInsert(values={"id": "s.id", "value": "s.value"}) \
.execute()
3. Write-then-rename pattern - запись во временную директорию, затем атомарный rename (для HDFS).
4. Deduplication перед записью:
df.dropDuplicates(["business_key"]).write.mode("append").parquet(path)
Главный принцип: хранить состояние последнего успешного запуска (водяной знак, версию снапшота) во внешнем хранилище и на его основе фильтровать данные при старте.
Что такое Compaction (VACUUM) в контексте Delta Lake?¶
Delta Lake накапливает историю: каждый INSERT/UPDATE/DELETE создаёт новые файлы, старые остаются (для Time Travel). Со временем накапливается много мелких файлов и устаревших данных.
Compaction (OPTIMIZE) - объединяет мелкие Parquet-файлы в крупные:
-- SQL: оптимизировать + опционально Z-order
OPTIMIZE events;
OPTIMIZE events ZORDER BY (user_id, event_date);
# Python API
from delta.tables import DeltaTable
dt = DeltaTable.forPath(spark, "s3://bucket/delta/events/")
dt.optimize().executeCompaction()
dt.optimize().executeZOrderBy("user_id", "event_date")
VACUUM - удаляет файлы, на которые больше нет ссылок в transaction log (старые версии):
-- Удалить файлы старше 7 дней (default retention)
VACUUM events;
-- Уменьшить retention (осторожно - сломает Time Travel за этот период)
VACUUM events RETAIN 168 HOURS;
-- Dry run - показать что будет удалено
VACUUM events DRY RUN;
Связь: OPTIMIZE уменьшает число файлов → меньше overhead при чтении. VACUUM освобождает дисковое пространство → дешевле хранение на S3/HDFS. Оба нужны в регулярном обслуживании таблицы.