Блок 3 - DataFrame API и SQL

20 вопросов: чтение форматов, join-типы, window functions, null-обработка, complex types и actions.

core

Как прочитать данные из 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. Оба нужны в регулярном обслуживании таблицы.