Локальная установка PySpark и первый job с нуля
Установка PySpark 4.x: Java 17, Python 3.10+, три пакета (pyspark / pyspark-connect / pyspark-client), классический режим и Spark Connect, первый SparkSession, spark-submit, Spark UI и типичные ошибки.
Этот урок - практический старт. Задача: установить 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¶
- Скачать JDK 17 с adoptium.net (Eclipse Temurin)
- Запустить установщик
.msi - Добавить системную переменную
JAVA_HOME = C:\Program Files\Eclipse Adoptium\jdk-17.x.x.x-hotspot - Добавить в
Path:%JAVA_HOME%\bin - Проверить в 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.0local[*]как master URL - все ядра машины, кластер не нужен- Spark UI на
:4040- доступен пока приложение работает
Важное новшество 4.x: Spark Connect - клиент-серверный режим подключения, который становится основным способом работы со Spark. Для изучения внутренностей мы пока используем классический режим - он даёт прямой доступ к DAG Scheduler, Executor-ам и Spark UI без посредников.