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, но размазанная по кластеру и вычисляемая лениво.
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 yet | Trigger actual computation |
select, filter, withColumn, join, groupBy, orderBy | show, collect, count, write, take |
| Just builds a logical plan (a DAG) | Submits a job, Spark runs the plan |
| Трансформации (ленивые) | Действия (немедленные) |
|---|---|
| Описывают новый DataFrame; пока ничего не выполняется | Запускают реальное вычисление |
select, filter, withColumn, join, groupBy, orderBy | show, 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
filter down to the data source, prune unused columns, and reorder steps. Eager execution couldn't do that.
Поскольку Spark видит всю цепочку до запуска, Catalyst может её оптимизировать: протолкнуть filter к источнику, отбросить неиспользуемые колонки, переставить шаги. При немедленном выполнении это было бы невозможно.
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
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"])
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")
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/"))
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 """)
8. Mini-glossary8. Мини-словарь
| Term | Plain meaning |
|---|---|
| SparkSession | Entry point / handle to the cluster (spark). |
| DataFrame | Distributed table with named, typed columns. |
| Transformation | Lazy op that builds a plan (select, filter, join). |
| Action | Eager op that runs the plan (show, count, write). |
| Catalyst | Spark's query optimizer (plans, pushdowns, reorders). |
| Partition | A chunk of the data processed by one task. |
| Driver | Coordinator process running your code & the plan. |
| Executor | Worker 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.Ответь про себя, потом нажми, чтобы перевернуть.
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 бы их схлопнул.