Spark Connect: gRPC thin client и удалённое подключение без локального Spark

Spark Connect - революция тонкого клиента: gRPC, Protocol Buffers, логический план как протокол, Arrow-транспорт результатов, ограничения и сценарии использования.

core

Архитектурный тупик классического 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 внутри).