Spark Connect: gRPC thin client и удалённое подключение без локального Spark
Spark Connect - революция тонкого клиента: gRPC, Protocol Buffers, логический план как протокол, Arrow-транспорт результатов, ограничения и сценарии использования.
Архитектурный тупик классического PySpark¶
До Spark 3.4 каждый PySpark-клиент требовал полного Spark-рантайма прямо на машине разработчика:
Проблемы этой архитектуры:
- Тяжёлый рантайм: установка
pysparkтянет JVM, Scala-библиотеки, сотни MB зависимостей - Конфликты версий: клиент и кластер должны быть одной версии Spark
- JVM на каждой машине разработчика: ноутбук аналитика - полноценный Driver
- Жёсткая привязка: если Python-процесс упал, Driver упал, все Task'и потеряны
- Сложность в микросервисах: FastAPI-сервис с Spark-запросами требует полного Spark-стека
Результат: Spark тяжело использовать из IDE, из Jupyter в облаке, из serverless-контейнеров.
Главная идея Spark Connect¶
Spark Connect (появился в Spark 3.4, GA в Spark 3.5) - это архитектура тонкого клиента: клиент знает только API (DataFrame, SQL), а весь Spark-рантайм живёт на сервере.
Клиент:
- Обычная Python-библиотека (
pip install pyspark-clientилиpip install pyspark) - Нет JVM, нет Py4J, нет локального Spark
- Строит логический план запроса локально
- Сериализует его в Protocol Buffers
- Отправляет на сервер через gRPC
- Получает результаты в формате Apache Arrow
Почему gRPC и Protocol Buffers¶
gRPC - не просто REST¶
gRPC - RPC-фреймворк от Google поверх HTTP/2. Ключевые свойства, которые важны для Spark Connect:
Protocol Buffers (protobuf) - язык описания структур данных. Spark Connect описывает всё - от DataFrame.filter() до GroupBy.agg() - в .proto-файлах:
// Упрощённый пример из spark/connect/proto/relations.proto
message Relation {
RelationCommon common = 1;
oneof rel_type {
Read read = 2;
Project project = 3;
Filter filter = 4;
Join join = 5;
Aggregate aggregate = 6;
Sort sort = 7;
// ...
}
}
message Filter {
Relation input = 1; // входное отношение (другая Relation)
Expression condition = 2; // условие фильтрации
}
Каждая операция df.filter(...) в Python создаёт protobuf-объект Filter, а не вызов Java-метода через Py4J.
Логический план как протокол передачи¶
Это самый важный концептуальный сдвиг Spark Connect. В классическом PySpark клиент управляет JVM-объектами через Py4J. В Spark Connect клиент описывает вычисление в виде логического плана.
Локальное построение дерева означает: вся цепочка .filter().groupBy().agg() строится как Python-объекты на клиенте. Никаких round-trip'ов к серверу на каждый шаг цепочки (в отличие от Py4J, где каждый .filter() - отдельный TCP-вызов).
Сравнение с классическим Py4J¶
| Аспект | Классический Py4J | Spark Connect |
|---|---|---|
| Клиентские зависимости | JVM + полный Spark + Py4J | Только Python-пакет |
| Число сетевых вызовов | По одному на каждый Transformation | Один на Action |
| Протокол | ASCII/binary через TCP socket | Protobuf через gRPC/HTTP2 |
| Transport результатов | Pickle (построчно) | Apache Arrow (батчи) |
| Версионная совместимость | Клиент = версия кластера | Клиент может отличаться |
| Изоляция клиентов | Один Driver = одно приложение | Сервер обслуживает много клиентов |
| RDD API | Доступен | Недоступен |
Архитектура Spark Connect Server¶
Spark Connect Server - это компонент внутри Spark Driver. Он запускается дополнительно к стандартному Driver'у и слушает gRPC-порт (по умолчанию :15002).
SparkConnectSessionManager управляет изоляцией клиентов. Каждый клиент получает свою SparkSession внутри сервера - отдельные временные таблицы, конфигурация, кэш. Падение одного клиента не затрагивает других.
Транспорт результатов: Apache Arrow¶
Вместо построчного Pickle (как в классических Python UDF) Spark Connect использует Apache Arrow для всех результатов:
Arrow позволяет:
- Передавать батчи (тысячи строк за один gRPC-фрейм)
- Создавать
pd.DataFrameиз Arrow-буфера без копирования (zero-copy) - Сохранять типы данных (int64, float64, string) без потерь при сериализации
Настройка и подключение¶
Запуск Spark Connect Server¶
# Запустить Spark Connect Server (Spark 3.4+)
./sbin/start-connect-server.sh \
--packages org.apache.spark:spark-connect_2.12:3.5.0 \
--master yarn \
--conf spark.executor.instances=10
# Или через spark-submit:
spark-submit \
--packages org.apache.spark:spark-connect_2.12:3.5.0 \
--class org.apache.spark.sql.connect.service.SparkConnectServer \
# (класс запускается автоматически через скрипт)
По умолчанию сервер слушает 0.0.0.0:15002.
Подключение клиента¶
# Классический PySpark (требует JVM + полный Spark):
spark = SparkSession.builder \
.master("yarn") \
.appName("MyJob") \
.getOrCreate()
# Spark Connect (только Python-пакет, без JVM):
spark = SparkSession.builder \
.remote("sc://spark-server.internal:15002") \
.getOrCreate()
# Строка подключения:
# sc://host:port[/;param=value;...]
# sc://spark-server:15002
# sc://spark-server:15002/;token=my-auth-token
# sc://spark-server:15002/;use_ssl=true
Установка клиента¶
# Минимальная установка для Spark Connect (без полного Spark):
pip install pyspark==3.5.0
# Или если нужен только клиент (Databricks):
pip install databricks-connect==14.0
Особенности UDF в Spark Connect¶
Python UDF через Spark Connect работают иначе, чем в классическом PySpark:
UDF-код сериализуется через CloudPickle на клиенте и передаётся внутри protobuf-плана на сервер, а оттуда - на Executor'ы. Механизм тот же, что в классическом PySpark, но маршрут другой: не Py4J → TaskScheduler, а gRPC → SparkConnectPlanner → TaskScheduler.
Добавление файлов и зависимостей¶
# В классическом PySpark:
spark.sparkContext.addPyFile("mymodule.py")
spark.sparkContext.addFile("config.json")
# В Spark Connect:
spark.addArtifacts("mymodule.py") # Python-модуль
spark.addArtifacts("config.json") # произвольный файл
spark.addArtifacts("mypackage.whl") # Python-пакет
# Файлы передаются на сервер через gRPC, а затем раздаются Executor'ам
Безопасность и мультитенантность¶
# Подключение с TLS и токеном аутентификации:
spark = SparkSession.builder \
.remote("sc://spark-cluster:15002/;use_ssl=true;token=Bearer_TOKEN_HERE") \
.getOrCreate()
Сервер может обслуживать множество клиентов одновременно. Каждая сессия изолирована: временные таблицы, конфигурация, закэшированные DataFrames одного клиента не видны другим.
Ограничения Spark Connect¶
Spark Connect - не полный эквивалент классического SparkSession. Некоторые возможности недоступны:
RDD API отсутствует принципиально: Spark Connect передаёт логические планы через protobuf, а RDD - это низкоуровневая абстракция, не описываемая в виде плана. Это намеренное ограничение - RDD API устаревает.
Сценарии использования¶
1. IDE разработка без установки Spark¶
# Раньше: нужен локальный Spark + JVM + JAVA_HOME
# Сейчас: только pip install pyspark
spark = SparkSession.builder \
.remote("sc://dev-spark-cluster:15002") \
.getOrCreate()
df = spark.read.parquet("s3://data-lake/events/")
df.filter(col("date") >= "2024-01-01").show()
2. Spark в микросервисе (FastAPI)¶
# Dockerfile: FROM python:3.11-slim
# pip install pyspark fastapi uvicorn
# Никакого JVM!
from fastapi import FastAPI
from pyspark.sql import SparkSession
app = FastAPI()
@app.on_event("startup")
async def startup():
global spark
spark = SparkSession.builder \
.remote("sc://spark-cluster:15002") \
.getOrCreate()
@app.get("/analytics/{user_id}")
async def get_user_stats(user_id: str):
result = spark.sql(f"""
SELECT region, count(*) as events
FROM events
WHERE user_id = '{user_id}'
GROUP BY region
""").collect()
return {"stats": [row.asDict() for row in result]}
3. Jupyter/Notebook в облаке¶
# JupyterHub без Spark-инсталляции на каждом pod'е
# Единый Spark Connect Server на кластере
# Все ноутбуки подключаются к нему
spark = SparkSession.builder \
.remote("sc://shared-spark:15002/;token=USER_TOKEN") \
.getOrCreate()
# Изоляция: каждый пользователь видит свои temp views
spark.sql("CREATE TEMP VIEW my_data AS SELECT ...")
4. Версионная совместимость¶
# Сервер: Spark 3.5 на Python 3.9
# Клиент: pyspark 3.5, Python 3.11 - совместимо!
# Правило совместимости:
# - Минорные версии: 3.5.x клиент ↔ 3.5.y сервер - любая комбинация
# - Мажорные версии: 3.x ↔ 4.x - обратная совместимость не гарантирована
Диагностика¶
Логи на стороне сервера¶
# Логи Spark Connect Server (внутри Driver лог-файлов):
grep "SparkConnectService\|ConnectSession" spark-driver.log
# Типичные записи:
# INFO SparkConnectService: Starting SparkConnect server on port 15002
# INFO ConnectSessionManager: New session: user=alice, sessionId=abc123
# INFO ConnectSessionManager: Session abc123 closed after 3600s
Диагностика на стороне клиента¶
import logging
# Включить gRPC-логирование
logging.basicConfig(level=logging.DEBUG)
logging.getLogger("grpc").setLevel(logging.DEBUG)
spark = SparkSession.builder.remote("sc://server:15002").getOrCreate()
# Проверить соединение
print(spark.version) # версия сервера
print(spark.sparkContext.master) # "spark://server:15002" (Connect URL)
Типичные ошибки¶
| Ошибка | Причина | Решение |
|---|---|---|
StatusCode.UNAVAILABLE |
Сервер недоступен / неверный порт | Проверить nc -zv host 15002 |
StatusCode.UNAUTHENTICATED |
Неверный токен | Проверить /;token=... в URL |
AnalysisException: UNSUPPORTED_FEATURE |
Вызов RDD API | Переписать через DataFrame API |
PicklingError в UDF |
Non-serializable замыкание | Инициализировать ресурсы внутри UDF |
Timeout при большом collect() |
gRPC timeout меньше времени выполнения | spark.conf.set("spark.connect.grpc.max.message.size", "128m") |
Практика¶
1. Запустить локальный Spark Connect Server¶
# Требования: Java 11+, Spark 3.5+ distribution
export SPARK_HOME=/opt/spark-3.5.0
# Запуск в local mode
$SPARK_HOME/sbin/start-connect-server.sh \
--master local[4] \
--packages org.apache.spark:spark-connect_2.12:3.5.0
# Проверить что слушает порт:
ss -tlnp | grep 15002
2. Подключиться и выполнить запрос¶
from pyspark.sql import SparkSession
from pyspark.sql.functions import col, count, avg
# Подключение через Spark Connect
spark = SparkSession.builder \
.remote("sc://localhost:15002") \
.getOrCreate()
print(f"Подключён к Spark {spark.version}")
# DataFrame API работает как обычно
df = spark.range(1_000_000) \
.selectExpr("id", "id % 100 AS bucket", "rand() * 100 AS value")
result = df.groupBy("bucket") \
.agg(count("*").alias("cnt"), avg("value").alias("avg_value")) \
.orderBy("bucket")
result.show(5)
# Результат приходит через gRPC как Arrow batches
3. UDF через Spark Connect¶
from pyspark.sql.functions import udf, pandas_udf
from pyspark.sql.types import StringType
import pandas as pd
# Python UDF - работает через CloudPickle → gRPC → Executor
@udf(StringType())
def classify(v):
if v > 75:
return "high"
elif v > 25:
return "medium"
return "low"
# Pandas UDF - работает через Arrow → gRPC → Executor
@pandas_udf(StringType())
def classify_batch(s: pd.Series) -> pd.Series:
return s.map(lambda v: "high" if v > 75 else "medium" if v > 25 else "low")
df.withColumn("class", classify("value")).show(3)
df.withColumn("class", classify_batch("value")).show(3)
4. Сравнение задержки: Connect vs классический¶
import time
# Spark Connect: цепочка transformations - локально, Action - один gRPC
spark_connect = SparkSession.builder.remote("sc://localhost:15002").getOrCreate()
t0 = time.time()
for _ in range(100):
# Каждый вызов - локальное построение дерева (без сети)
_ = spark_connect.range(100).filter(col("id") > 50)
# Ни одного сетевого вызова - всё локально!
print(f"Spark Connect (100 transformations, no action): {(time.time()-t0)*1000:.1f}ms")
# Один Action - один gRPC
t0 = time.time()
spark_connect.range(100).filter(col("id") > 50).count()
print(f"Spark Connect (1 action): {(time.time()-t0)*1000:.1f}ms")
# В классическом Py4J каждый .filter(), .select() - это TCP round-trip
Будущее: Spark Connect как стандарт¶
Spark Connect - не просто новая фича, а изменение парадигмы:
Тенденции:
- Databricks Connect (основан на Spark Connect) - стандарт разработки на Databricks
- Spark Serverless (AWS Glue, Azure Synapse) - уже используют similar thin-client модель
- Multi-language clients - Spark Connect .proto файлы позволяют писать клиентов на Go, Rust, TypeScript
- IDE интеграции - VS Code и IntelliJ уже имеют плагины для Spark Connect
Резюме¶
Spark Connect решает главную проблему классического PySpark - жёсткую связь клиента с Spark-рантаймом:
- gRPC + Protocol Buffers вместо Py4J TCP-сокета - языконейтральный бинарный протокол
- Логический план как payload - клиент строит дерево локально, отправляет одним запросом
- Apache Arrow для результатов - батчевый колоночный транспорт вместо построчного Pickle
- Изоляция сессий - один сервер, много клиентов, каждый в своей SparkSession
- Без JVM на клиенте -
pip install pysparkдостаточно для работы с удалённым кластером
Ограничения: нет RDD API (намеренно), нет прямого доступа к JVM, UDF передаются через protobuf (CloudPickle внутри).