Локальная установка PySpark и первый job с нуля

Установка PySpark 4.x: Java 17, Python 3.10+, три пакета (pyspark / pyspark-connect / pyspark-client), классический режим и Spark Connect, первый SparkSession, spark-submit, Spark UI и типичные ошибки.

core

Этот урок - практический старт. Задача: установить PySpark 4.x, написать первый job и убедиться, что всё работает до того, как двигаться дальше. Без локальной среды сложно проверить интуицию о DAG, трансформациях и shuffle на реальном коде.


Почему PySpark требует Java

PySpark - это Python-обёртка над Apache Spark, который работает на JVM (Java Virtual Machine). Когда вы вызываете df.groupBy().count(), Python общается с JVM через Py4J - RPC-мост между Python и Java. Вся реальная работа (планирование, выполнение, shuffle) происходит в JVM.

В local mode (master("local[*]")) Driver и все «Executor-ы» - это потоки одного JVM-процесса на вашей машине. Кластер не нужен.


Требования для PySpark 4.x

Компонент Версия
Java (JDK) 17 или выше (обязательно)
Python 3.10 или выше
RAM 8 GB рекомендуется, 4 GB минимум

PySpark 4.x поддерживает только Java 17+. Python 3.9 и ниже не поддерживается.


Три пакета PySpark 4.x: какой выбрать

В PySpark 4.x появились три отдельных pip-пакета с разным назначением:

pyspark (классический режим) - полный Spark Runtime. SparkSession запускается прямо в вашем Python-процессе, всё работает локально. Требует Java 17+.

pyspark-connect - Spark Connect как режим по умолчанию. При SparkSession.builder.getOrCreate() автоматически запускает локальный Spark Connect сервер. Поддерживает и master=local[*], и remote=sc://.... Требует Java 17+.

pyspark-client - чистый Python-клиент без Java. Подключается только к удалённому Spark Connect серверу через sc://hostname:port. Подходит для сред без JVM.

Для локального обучения используем pyspark[sql] - классический режим с минимальными зависимостями.


Шаг 1: Установка Java 17

Linux (Ubuntu / Debian)

sudo apt update
sudo apt install -y openjdk-17-jdk

# Проверка
java -version
# openjdk version "17.x.x" ...

# Установить JAVA_HOME
echo 'export JAVA_HOME=$(dirname $(dirname $(readlink -f $(which java))))' >> ~/.bashrc
echo 'export PATH=$JAVA_HOME/bin:$PATH' >> ~/.bashrc
source ~/.bashrc

echo $JAVA_HOME
# /usr/lib/jvm/java-17-openjdk-amd64

macOS

# Через Homebrew
brew install openjdk@17

# Добавить в PATH (для zsh)
echo 'export JAVA_HOME=$(brew --prefix openjdk@17)' >> ~/.zshrc
echo 'export PATH=$JAVA_HOME/bin:$PATH' >> ~/.zshrc
source ~/.zshrc

java -version

Windows

  1. Скачать JDK 17 с adoptium.net (Eclipse Temurin)
  2. Запустить установщик .msi
  3. Добавить системную переменную JAVA_HOME = C:\Program Files\Eclipse Adoptium\jdk-17.x.x.x-hotspot
  4. Добавить в Path: %JAVA_HOME%\bin
  5. Проверить в PowerShell: java -version

Windows: winutils.exe

Spark на Windows требует winutils.exe. Без него возникает ошибка ERROR Shell: Failed to locate the winutils binary:

# Скачать winutils для Hadoop 3.x с https://github.com/cdarlint/winutils
mkdir C:\hadoop\bin
# Скопировать winutils.exe и hadoop.dll в C:\hadoop\bin

[System.Environment]::SetEnvironmentVariable("HADOOP_HOME", "C:\hadoop", "Machine")

Шаг 2: Python-окружение

Официальная документация PySpark не рекомендует смешивать pip и conda в одном окружении - это приводит к конфликтам зависимостей. Выберите один инструмент и придерживайтесь его.

Вариант A: venv (рекомендуется для pip)

# Создать окружение (Python 3.10+)
python3 -m venv spark-env

# Активировать
source spark-env/bin/activate        # Linux / macOS
# spark-env\Scripts\activate.bat     # Windows CMD
# spark-env\Scripts\Activate.ps1     # Windows PowerShell

# Проверить
which python  # → spark-env/bin/python
python --version  # → Python 3.10.x или выше

Вариант B: conda (устанавливать PySpark только через conda)

conda create -n spark-env python=3.11 -y
conda activate spark-env

# Устанавливать через conda, не через pip
conda install -c conda-forge pyspark

Conda-пакеты могут выходить с задержкой относительно PyPI-релизов. Если нужна последняя версия - используйте venv + pip.


Шаг 3: Установка PySpark

# Рекомендуемый вариант: классический режим с SQL-зависимостями
pip install "pyspark[sql]"

# Дополнительно для Data Engineer стека
pip install "pyspark[sql]" delta-spark

# Spark Connect режим (Java всё равно нужна)
pip install "pyspark[connect]"

# Все компоненты сразу
pip install "pyspark[sql,connect,pandas_on_spark]"

# Проверка
python -c "import pyspark; print(pyspark.__version__)"
# 4.0.0 (или 4.1.x в зависимости от актуальной версии)

Что включает pyspark[sql]:

Пакет Версия Назначение
py4j ≥ 0.10.9.9 Python ↔ JVM мост
pandas ≥ 2.2.0 pandas UDF, toPandas()
pyarrow ≥ 15.0.0 Arrow-сериализация, Parquet

pip install pyspark скачивает и распаковывает полный дистрибутив Spark (~300 MB). SPARK_HOME устанавливается автоматически внутри пакета - отдельно скачивать Spark не нужно.


Шаг 4: Проверка через pyspark shell

pyspark --master local[*]

Вы увидите ASCII-баннер Spark и приглашение >>>. Это интерактивная оболочка с готовым объектом spark:

# Внутри pyspark shell
spark.range(10).show()
# +---+
# | id|
# +---+
# |  0|
# |  1|
# ...
# |  9|
# +---+

spark.version
# '4.0.0'

# Открыть Spark UI в браузере: http://localhost:4040

Первый SparkSession: разбор по строкам

from pyspark.sql import SparkSession
from pyspark.sql.functions import col, upper, count

# 1. Создать SparkSession
spark = (
    SparkSession.builder
    .master("local[*]")          # все ядра машины как Executor-ы
    .appName("FirstJob")         # имя в Spark UI
    .config("spark.driver.memory", "2g")
    .config("spark.sql.shuffle.partitions", "4")  # меньше для локального
    .getOrCreate()
)

# 2. Создать DataFrame из данных в памяти
data = [
    (1, "Alice",  "Moscow",  5000),
    (2, "Bob",    "Berlin",  7200),
    (3, "Carol",  "Moscow",  4800),
    (4, "Dave",   "Paris",   9100),
    (5, "Eve",    "Berlin",  6300),
    (6, "Frank",  "Moscow",  5500),
]
columns = ["id", "name", "city", "salary"]

df = spark.createDataFrame(data, columns)

# Альтернатива: с явной схемой через Row
from pyspark.sql import Row
df2 = spark.createDataFrame([
    Row(id=1, name="Alice", city="Moscow", salary=5000),
    Row(id=2, name="Bob",   city="Berlin", salary=7200),
])

# 3. Трансформации (lazy - ничего не выполняется)
result = (
    df.filter(col("salary") > 5000)
      .withColumn("name_upper", upper(col("name")))
      .groupBy("city")
      .agg(count("*").alias("employees"))
)

# 4. Action - запускает выполнение
result.show()
# +------+---------+
# |  city|employees|
# +------+---------+
# |Berlin|        2|
# | Paris|        1|
# +------+---------+

# 5. Схема и описание
df.printSchema()
# root
#  |-- id: long (nullable = true)
#  |-- name: string (nullable = true)
#  |-- city: string (nullable = true)
#  |-- salary: long (nullable = true)

df.select("salary", "city").describe().show()

# 6. Остановить сессию
spark.stop()

Что происходит при result.show():


Чтение и запись файлов

# CSV с заголовком
df_csv = spark.read.csv("data/users.csv", header=True, inferSchema=True)

# JSON
df_json = spark.read.json("data/events.json")

# Parquet (рекомендуемый формат для production)
df_parquet = spark.read.parquet("data/orders/")

# С явной схемой - надёжнее inferSchema
from pyspark.sql.types import StructType, StructField, IntegerType, StringType

schema = StructType([
    StructField("id",   IntegerType(), nullable=False),
    StructField("name", StringType(),  nullable=True),
])
df_typed = spark.read.csv("data/users.csv", header=True, schema=schema)

# Запись
df.write.mode("overwrite").parquet("output/result/")
df.write.mode("append").partitionBy("city").parquet("output/partitioned/")
df.write.csv("output/result.csv", header=True)

Работа с SQL

# Зарегистрировать DataFrame как временное представление
df.createOrReplaceTempView("employees")

# Запросить через SQL
result = spark.sql("""
    SELECT city, COUNT(*) AS cnt
    FROM employees
    WHERE salary > 5000
    GROUP BY city
    ORDER BY cnt DESC
""")
result.show()

# selectExpr - SQL-выражения прямо в DataFrame API
df.selectExpr("city", "salary * 1.1 AS salary_with_bonus").show()

Запуск через spark-submit

Сохраните код в файл и запустите как production-like job:

# word_count.py
from pyspark.sql import SparkSession
from pyspark.sql.functions import explode, split, lower, col, count


def main():
    spark = (
        SparkSession.builder
        .master("local[*]")
        .appName("WordCount")
        .getOrCreate()
    )

    lines = spark.createDataFrame([
        ("Apache Spark is a unified analytics engine",),
        ("Spark provides an interface for programming clusters",),
        ("PySpark is the Python API for Apache Spark",),
    ], ["line"])

    words = (
        lines
        .select(explode(split(lower(col("line")), r"\s+")).alias("word"))
        .filter(col("word") != "")
        .groupBy("word")
        .agg(count("*").alias("cnt"))
        .orderBy(col("cnt").desc())
    )

    words.show(20, truncate=False)
    spark.stop()


if __name__ == "__main__":
    main()
spark-submit word_count.py

Вывод:

+----------+---+
|word      |cnt|
+----------+---+
|spark     |3  |
|apache    |2  |
|is        |2  |
|for       |2  |
...


Spark Connect: новый режим в PySpark 4.x

Spark Connect - клиент-серверная архитектура, появившаяся в Spark 3.4 и ставшая приоритетным режимом в 4.x. Клиент (ваш Python-код) общается с Spark-сервером по gRPC, а не напрямую через JVM.

Преимущества Spark Connect:

  • Клиент изолирован от сервера - версия Python/библиотек на клиенте не влияет на сервер
  • pyspark-client не требует Java на машине разработчика
  • Поддерживает подключение к удалённым кластерам из ноутбука
  • Более стабильный API для интерактивной работы

Запуск локального Spark Connect сервера

# Установить с Connect-зависимостями
pip install "pyspark[connect]"

# Запустить локальный сервер (порт 15002)
$SPARK_HOME/sbin/start-connect-server.sh --master local[*]
# или, если SPARK_HOME настроен через pip:
python -m pyspark.sql.connect.server --master local[*]

Подключение к серверу

from pyspark.sql import SparkSession

# Подключиться к локальному серверу
spark = SparkSession.builder.remote("sc://localhost").getOrCreate()

# Работа идентична классическому режиму
df = spark.createDataFrame([(1, "Alice"), (2, "Bob")], ["id", "name"])
df.show()

spark.stop()

Автоматический Connect через pyspark-connect

pip install pyspark-connect
from pyspark.sql import SparkSession

# Автоматически запускает локальный Connect-сервер
spark = SparkSession.builder.getOrCreate()

Для целей этого курса мы работаем в классическом режиме (master("local[*]")): он проще в настройке и позволяет увидеть все внутренности Spark напрямую. Spark Connect изучим в разделе по production-деплойменту.


Spark UI: что смотреть

Пока Spark-приложение работает, откройте http://localhost:4040:

Для WordCount вы увидите:

  • Jobs - 1 job (от show())
  • Stages - 2 stage: explode + filter (narrow) и groupBy + orderBy (2 shuffle)
  • Environment - все параметры конфигурации, версии Java и Spark

После завершения приложения UI недоступен. Для анализа завершённых job нужен History Server (настроим в модуле по мониторингу).


Настройка Jupyter Notebook

pip install "pyspark[sql]" jupyter notebook ipykernel

# Запустить обычным способом
jupyter notebook

В ноутбуке PySpark используется как обычная библиотека:

import sys
import os
from pyspark.sql import SparkSession
from pyspark.sql.functions import col

# Убедиться, что Driver и воркеры используют одинаковый Python
os.environ["PYSPARK_PYTHON"] = sys.executable
os.environ["PYSPARK_DRIVER_PYTHON"] = sys.executable

spark = (
    SparkSession.builder
    .master("local[4]")
    .appName("Notebook")
    .config("spark.driver.memory", "4g")
    .config("spark.sql.shuffle.partitions", "4")
    .getOrCreate()
)

spark.range(100) \
    .groupBy((col("id") % 10).alias("bucket")) \
    .count() \
    .show()

Важно: один SparkSession на всё время жизни ядра ноутбука. Для смены конфигурации - Kernel → Restart.

Устаревший паттерн PYSPARK_DRIVER_PYTHON="jupyter" (из старых руководств) использовался для запуска Jupyter из pyspark shell - он не нужен при обычной работе с jupyter notebook.


Типичные проблемы и решения

Java не найдена или версия старше 17

Exception in thread "main" java.lang.UnsupportedClassVersionError
# или
JAVA_HOME is not set
# или
Unsupported class file major version 55 (Java 11 → ошибка в PySpark 4.x)
# Проверить
java -version
echo $JAVA_HOME

# Исправить: установить Java 17 и обновить JAVA_HOME
sudo apt install openjdk-17-jdk         # Linux
brew install openjdk@17                  # macOS
export JAVA_HOME=/usr/lib/jvm/java-17-openjdk-amd64

PySpark 4.x требует Java 17 как минимум. Java 11 не поддерживается.

Python ниже 3.10

ERROR: Package 'pyspark' requires a different Python:
3.9.x not in '>=3.10'
# Проверить версию
python --version

# Создать окружение с правильной версией
python3.11 -m venv spark-env
source spark-env/bin/activate

Порт 4040 занят

Spark автоматически пробует 4041, 4042... Для фиксированного порта:

.config("spark.ui.port", "4050")

Python worker failed to connect

Exception: Python worker failed to connect back

Driver и Executor используют разные Python-интерпретаторы:

import sys
os.environ["PYSPARK_PYTHON"] = sys.executable
os.environ["PYSPARK_DRIVER_PYTHON"] = sys.executable

OutOfMemoryError

java.lang.OutOfMemoryError: Java heap space
# Увеличить память Driver (local mode = Driver + Executor в одном JVM)
.config("spark.driver.memory", "4g")

# Уменьшить данные при разработке
df.sample(0.01)  # 1% данных для отладки

Spark слишком много пишет в логи

spark.sparkContext.setLogLevel("WARN")  # ERROR / WARN / INFO / DEBUG

Чек-лист готовности

# 1. Java 17+ работает
java -version
# → openjdk version "17.x.x"

# 2. Python 3.10+ активен
python --version
# → Python 3.11.x

# 3. PySpark 4.x установлен
python -c "import pyspark; print(pyspark.__version__)"
# → 4.x.x

# 4. Простейший job выполняется без ошибок
python - <<'EOF'
from pyspark.sql import SparkSession
spark = SparkSession.builder.master("local[*]").appName("test").getOrCreate()
result = spark.range(10).count()
assert result == 10, f"Expected 10, got {result}"
print(f"OK: spark.range(10).count() = {result}")
print(f"Spark version: {spark.version}")
spark.stop()
EOF
# → OK: spark.range(10).count() = 10
# → Spark version: 4.x.x

# 5. spark-submit работает
spark-submit --master local[*] --conf spark.sql.shuffle.partitions=4 word_count.py

Если все 5 проверок прошли - можно двигаться дальше.


Итог

Локальная установка PySpark 4.x требует:

  • Java JDK 17+ - обязательно, Java 11 не поддерживается
  • Python 3.10+ - обязательно
  • pip install "pyspark[sql]" - включает pandas ≥ 2.2 и pyarrow ≥ 15.0
  • local[*] как master URL - все ядра машины, кластер не нужен
  • Spark UI на :4040 - доступен пока приложение работает

Важное новшество 4.x: Spark Connect - клиент-серверный режим подключения, который становится основным способом работы со Spark. Для изучения внутренностей мы пока используем классический режим - он даёт прямой доступ к DAG Scheduler, Executor-ам и Spark UI без посредников.