PySpark EssentialsPySpark: основы

The DataFrame API you actually write at work: lazy evaluation, transformations vs actions, select/filter/join/groupBy, schemas, read/write. Code-first. When comfortable, switch to Performance. DataFrame API, который реально пишешь на работе: ленивые вычисления, трансформации vs действия, select/filter/join/groupBy, схемы, чтение/запись. Сначала код. Когда станет комфортно, переключайся на Performance.
🟢 EssentialsОсновы 🔴 PerformanceПроизводительность 📘 Spark theoryТеория Spark

1. What is PySpark?1. Что такое PySpark?

PySpark is the Python API for Apache Spark, a distributed engine for processing data that's too big for one machine. You write Python; Spark splits the work across many machines (a cluster) and runs it in parallel. The main thing you touch is the DataFrame: a distributed table with named, typed columns — like a pandas DataFrame, but spread over the cluster and computed lazily.

PySpark — это Python-API для Apache Spark, движка распределённой обработки данных, которые не помещаются на одной машине. Ты пишешь Python; Spark разбивает работу по множеству машин (кластер) и выполняет параллельно. Главное, с чем работаешь — DataFrame: распределённая таблица с именованными типизированными колонками — как pandas DataFrame, но размазанная по кластеру и вычисляемая лениво.

AnalogyАналогия Pandas is one cook in one kitchen. PySpark is a head chef writing a recipe, then handing copies to 100 cooks who each prep a slice of the ingredients at the same time. You describe what you want; Spark figures out how to divide and run it. Pandas — один повар на одной кухне. PySpark — шеф, пишущий рецепт, и раздающий копии 100 поварам, каждый из которых одновременно готовит свою часть продуктов. Ты описываешь что хочешь; Spark решает, как разбить и выполнить.
One lineОдной строкой PySpark = Python + DataFrames over a cluster. You write declarative transformations; Spark plans and runs them in parallel. PySpark = Python + DataFrame поверх кластера. Ты пишешь декларативные трансформации; Spark планирует и выполняет их параллельно.
DataFrame vs RDDDataFrame vs RDD RDD is the old low-level API (rows are opaque objects). Almost always use DataFrame / Spark SQL: it goes through the Catalyst optimizer, so it's faster and shorter. Only drop to RDD for rare custom logic. RDD — старый низкоуровневый API (строки — непрозрачные объекты). Почти всегда используй DataFrame / Spark SQL: он проходит через оптимизатор Catalyst, поэтому быстрее и короче. Спускайся к RDD только для редкой кастомной логики.

2. Lazy evaluation: transformations vs actions2. Ленивые вычисления: трансформации vs действия

This is the single most important PySpark idea, and a near-guaranteed interview question. Operations come in two kinds:

Это самая важная идея PySpark и почти гарантированный вопрос на интервью. Операции бывают двух видов:

Transformations (lazy)Actions (eager)
Describe a new DataFrame; nothing runs yetTrigger actual computation
select, filter, withColumn, join, groupBy, orderByshow, collect, count, write, take
Just builds a logical plan (a DAG)Submits a job, Spark runs the plan
Трансформации (ленивые)Действия (немедленные)
Описывают новый DataFrame; пока ничего не выполняетсяЗапускают реальное вычисление
select, filter, withColumn, join, groupBy, orderByshow, collect, count, write, take
Просто строит логический план (DAG)Отправляет job, Spark выполняет план
# nothing runs here — just building a plan
df2 = (df
    .filter(df.country == "DE")
    .withColumn("amount_eur", df.amount * 1.0)
    .select("user_id", "amount_eur"))

df2.show()   # ACTION → now Spark optimizes & executes the whole chain
Why lazy is goodПочему лень — это хорошо Because Spark sees the whole chain before running, Catalyst can optimize it: push your filter down to the data source, prune unused columns, and reorder steps. Eager execution couldn't do that. Поскольку Spark видит всю цепочку до запуска, Catalyst может её оптимизировать: протолкнуть filter к источнику, отбросить неиспользуемые колонки, переставить шаги. При немедленном выполнении это было бы невозможно.
Interview trapЛовушка на интервью collect() pulls ALL rows to the driver's memory — fine for small results, an OOM crash on big data. Prefer show(20), take(n), or write to storage. Same for toPandas(). collect() тянет ВСЕ строки в память драйвера — норм для маленьких результатов, OOM-крах на больших данных. Лучше show(20), take(n) или запись в хранилище. То же про toPandas().

3. Getting going: SparkSession & schemas3. Старт: SparkSession и схемы

The SparkSession is your entry point — the handle to the cluster. Everything starts from spark.

SparkSession — твоя точка входа, ручка к кластеру. Всё начинается с spark.

from pyspark.sql import SparkSession
from pyspark.sql import functions as F
from pyspark.sql.types import StructType, StructField, StringType, IntegerType

spark = (SparkSession.builder
    .appName("demo")
    .getOrCreate())

# create from Python data (handy for tests)
df = spark.createDataFrame(
    [(1, "Amal", "DE"), (2, "Tanya", "UK")],
    ["id", "name", "country"])

df.printSchema()   # inspect column names & types
df.show()          # pretty-print first 20 rows
Always declare a schema for production readsВ проде всегда задавай схему при чтении inferSchema=True makes Spark scan the data twice and can guess types wrong. Pass an explicit StructType — faster and safe. inferSchema=True заставляет Spark сканировать данные дважды и может угадать типы неверно. Передавай явный StructType — быстрее и надёжнее.
schema = StructType([
    StructField("id", IntegerType(), False),
    StructField("name", StringType(), True),
    StructField("country", StringType(), True),
])
df = spark.read.schema(schema).csv("s3://bucket/users.csv", header=True)

4. Core operations (the daily 80%)4. Основные операции (ежедневные 80%)

Three ways to reference a column: "name" (string), df.name (attribute), F.col("name") (function). Use F.col when you need expressions.

Три способа сослаться на колонку: "name" (строка), df.name (атрибут), F.col("name") (функция). Используй F.col, когда нужны выражения.

# select & rename
df.select("id", F.col("name").alias("user_name"))

# filter (where == filter)
df.filter((F.col("country") == "DE") & (F.col("id") > 0))

# add / replace a column
df.withColumn("is_de", F.col("country") == "DE")

# conditional logic
df.withColumn("tier",
    F.when(F.col("amount") > 100, "hi")
     .when(F.col("amount") > 10, "mid")
     .otherwise("lo"))

# drop, distinct, dedup, sort, limit
df.drop("tmp").dropDuplicates(["id"]).orderBy(F.col("amount").desc()).limit(100)

# nulls
df.na.fill({"amount": 0}).na.drop(subset=["id"])
Use F functions, not Python loopsИспользуй функции F, а не циклы Python Built-in pyspark.sql.functions run inside the JVM engine and are optimized. Iterating rows in Python (or a plain UDF) is far slower. Reach for F.* first. Встроенные pyspark.sql.functions выполняются внутри JVM-движка и оптимизированы. Перебор строк в Python (или обычный UDF) намного медленнее. Сначала тянись к F.*.

5. Group by, aggregate & join5. Группировка, агрегация и join

# aggregation
(df.groupBy("country")
   .agg(
       F.count("*").alias("n"),
       F.sum("amount").alias("total"),
       F.avg("amount").alias("avg_amt"),
       F.countDistinct("user_id").alias("users"),
   ))

Joins — know the how values and the duplicate-column gotcha:

Join — знай значения how и подвох с дублирующимися колонками:

# join on a shared column name → no duplicate "country" column
orders.join(users, on="user_id", how="left")

# join on a condition → both id columns survive (ambiguous!)
orders.join(users, orders.user_id == users.id, "inner")

# how: inner | left | right | outer | left_semi | left_anti
# left_anti = rows in orders with NO match in users (great for "missing")
Window functions ≠ groupByОконные функции ≠ groupBy groupBy collapses rows; a window keeps every row and adds a computed column (running total, rank, row_number). Classic "dedup keeping latest": groupBy схлопывает строки; окно сохраняет каждую строку и добавляет вычисленную колонку (накопит. сумма, ранг, row_number). Классика «дедуп с сохранением последней»:
from pyspark.sql.window import Window
w = Window.partitionBy("user_id").orderBy(F.col("ts").desc())
latest = (df
    .withColumn("rn", F.row_number().over(w))
    .filter("rn = 1")
    .drop("rn"))

6. Reading & writing6. Чтение и запись

# read
df = spark.read.parquet("s3://bucket/events/")
df = spark.read.option("header", True).csv("path.csv")
df = spark.read.json("path.json")

# write — note mode & partitioning
(df.write
   .mode("overwrite")            # overwrite | append | ignore | error
   .partitionBy("ds")            # physical folder per partition value
   .parquet("s3://bucket/out/"))
Prefer columnar formatsПредпочитай колоночные форматы Parquet (or ORC) over CSV/JSON: compressed, typed, and supports column pruning + predicate pushdown — Spark reads only the columns and row-groups it needs. CSV reads everything. Parquet (или ORC) вместо CSV/JSON: сжатый, типизированный, поддерживает отсечение колонок + проталкивание предикатов — Spark читает только нужные колонки и row-group'ы. CSV читает всё.
partitionBy is for layout, not speed-everywherepartitionBy — про раскладку, не «ускоряет всё» Partition by a low-cardinality column you filter on often (e.g. date ds). Partitioning by a high-cardinality column (user_id) creates millions of tiny files — the "small files problem". Партиционируй по колонке с низкой кардинальностью, по которой часто фильтруешь (например дата ds). Партиционирование по высококардинальной колонке (user_id) создаёт миллионы крошечных файлов — «проблема мелких файлов».

7. DataFrame API vs Spark SQL7. DataFrame API vs Spark SQL

They're the same engine — pick whichever reads cleaner. Both go through Catalyst, so performance is identical.

Это один движок — выбирай, что читается чище. Оба идут через Catalyst, поэтому производительность одинакова.

# register a temp view, then write plain SQL
df.createOrReplaceTempView("orders")
top = spark.sql("""
    SELECT country, SUM(amount) AS total
    FROM orders
    WHERE amount > 0
    GROUP BY country
    ORDER BY total DESC
""")
Rule of thumbЭмпирическое правило Long, set-based logic → SQL reads better. Programmatic / parameterized / chained pipelines → DataFrame API. Mix freely. Длинная логика над множествами → SQL читается лучше. Программные / параметризованные / цепочечные пайплайны → DataFrame API. Смешивай свободно.

8. Mini-glossary8. Мини-словарь

TermPlain meaning
SparkSessionEntry point / handle to the cluster (spark).
DataFrameDistributed table with named, typed columns.
TransformationLazy op that builds a plan (select, filter, join).
ActionEager op that runs the plan (show, count, write).
CatalystSpark's query optimizer (plans, pushdowns, reorders).
PartitionA chunk of the data processed by one task.
DriverCoordinator process running your code & the plan.
ExecutorWorker process that runs tasks on partitions.
ТерминПростой смысл
SparkSessionТочка входа / ручка к кластеру (spark).
DataFrameРаспределённая таблица с именованными типизированными колонками.
ТрансформацияЛенивая операция, строящая план (select, filter, join).
Действие (action)Немедленная операция, запускающая план (show, count, write).
CatalystОптимизатор запросов Spark (планы, pushdown, перестановки).
ПартицияКусок данных, обрабатываемый одной задачей (task).
DriverПроцесс-координатор, выполняющий твой код и план.
ExecutorРабочий процесс, выполняющий задачи на партициях.

9. Quick self-check9. Быстрая самопроверка

Answer in your head, then tap to flip.Ответь про себя, потом нажми, чтобы перевернуть.

Transformation vs action?Трансформация vs действие?
tapнажми
Transformations (select/filter/join) are lazy — they only build a plan. Actions (show/count/write/collect) are eager — they trigger execution of the whole chain.Трансформации (select/filter/join) ленивые — лишь строят план. Действия (show/count/write/collect) немедленные — запускают выполнение всей цепочки.
Why avoid collect() on big data?Почему collect() опасен на больших данных?
tapнажми
It pulls every row to the driver's single-machine memory → OOM. Use show/take/limit or write to storage instead.Он тянет все строки в одномашинную память драйвера → OOM. Используй show/take/limit или запись в хранилище.
DataFrame vs RDD — which & why?DataFrame vs RDD — что и почему?
tapнажми
Use DataFrame: it has a schema and goes through Catalyst, so it's optimized and shorter. RDD is low-level/opaque — only for rare custom logic.Используй DataFrame: у него есть схема и он идёт через Catalyst, поэтому оптимизирован и короче. RDD низкоуровневый/непрозрачный — только для редкой кастомной логики.
Keep the latest row per user — how?Оставить последнюю строку на пользователя — как?
tapнажми
Window: row_number().over(partitionBy(user).orderBy(ts.desc())), then filter rn == 1. Window keeps rows; groupBy would collapse them.Окно: row_number().over(partitionBy(user).orderBy(ts.desc())), потом filter rn == 1. Окно сохраняет строки; groupBy бы их схлопнул.
Why Parquet over CSV?Почему Parquet вместо CSV?
tapнажми
Columnar + compressed + typed; supports column pruning and predicate pushdown so Spark reads only needed columns/row-groups. CSV reads everything and has no types.Колоночный + сжатый + типизированный; поддерживает отсечение колонок и pushdown предикатов, поэтому Spark читает только нужное. CSV читает всё и без типов.
Is Spark SQL slower than the DataFrame API?Spark SQL медленнее DataFrame API?
tapнажми
No — same engine, both compile through Catalyst to the same plan. Pick by readability.Нет — один движок, оба компилируются через Catalyst в один план. Выбирай по читаемости.