SPYT Connect
SPYT Connect — это механизм удалённого подключения к Spark-драйверу, построенный на базе протокола Spark Connect. С его помощью можно выполнять запросы Spark SQL через Query Tracker в YTsaurus. Также можно работать с данными напрямую из Python-кода через Spark Connect API, не устанавливая JVM на клиентской стороне.
Примечание
Механизм пришёл на смену Livy начиная со SPYT 2.10.0 и Query Tracker 0.4.
Когда возможны задержки запросов
SPYT Connect запускает Spark-драйвер по требованию. Если драйвер в данный момент не активен, его нужно сначала запустить — это занимает время. Драйвер может быть не активен в трёх случаях:
- Первое обращение к SPYT Connect.
Вы только начали работу, и драйвер запускается с нуля. - После простоя.
Чтобы не тратить ресурсы впустую, драйвер автоматически останавливается после 10 минут бездействия. Следующий запрос снова инициирует его запуск. Таймаут простоя настраивается через параметрspark.ytsaurus.connect.idle.timeoutв конфигурационных параметрах SPYT. - После изменения настроек.
Если вы поменяли конфигурацию ресурсов для запроса (например, количество ядер), старый драйвер останавливается, а под новые настройки запускается новый.
То, в какой момент вы заметите задержку инициализации, зависит от того, как вы работаете со SPYT Connect:
- В Query Tracker (UI или API) запуск сессии и отправка запроса в SPYT Connect происходят в рамках одного QT-запроса. Поэтому при запуске в интерфейсе или отправке запроса через API кажется, что сам запрос выполняется долго. Все последующие запросы будут отрабатывать быстро.
- В Python (Spark Connect API) вы управляете запуском явно. Задержка произойдёт ровно в тот момент, когда вы вызываете функцию
start_connect_server(ожидание готовности). Сами вычисления и операции с DataFrame будут стартовать без задержек.
Режимы запуска
SPYT Connect работает в двух режимах — они определяют, как запускается Spark-приложение. Выбор режима влияет на конфигурацию и код во всех способах подключения.
|
Режим |
Когда подходит |
|
Нет выделенного кластера; Spark-приложение запускается по требованию под каждый запрос |
|
|
Кластер уже запущен; SPYT Connect подключается к нему |
Способы подключения
|
Способ |
Когда подходит |
|
Подходит аналитикам и всем, кто работает с данными через интерфейс YTsaurus |
|
|
Для автоматизации SQL-запросов из Python |
|
|
Для тех, кто хочет использовать DataFrame API или управлять жизненным циклом драйвера вручную |
Через UI Query Tracker
Чтобы выполнить SQL-запрос с использованием SPYT Connect:
- Откройте вкладку Queries в интерфейсе YTsaurus.
- В списке движков выберите SPYT.
- Введите SQL-запрос.
- В поле Settings укажите конфигурацию в формате JSON.
- Нажмите Run и дождитесь результата.

Для работы с внутренним Spark-кластером добавьте discovery_path — путь к запущенному кластеру. Кластер должен работать на SPYT 2.9.0 или выше:
{
"discovery_path": "//home/spark/my-cluster"
}

Через Query Tracker API
В примере ниже показано, как отправить SQL-запрос через API и прочитать результат:
from yt.wrapper import YtClient, start_query, get_query_result, read_query_result
client = YtClient(proxy="<cluster-proxy>", token="<your-token>")
settings = {
"cluster": "<cluster-name>",
"spark_conf": {
"spark.cores.max": "4" # Spark-native параметр: максимальное число ядер для всего приложения
}
}
# Для внутреннего кластера добавьте discovery_path:
# settings["discovery_path"] = "//home/spark/my-cluster"
query_id = start_query(
"spyt",
"SELECT * FROM yt.`//home/my-table`",
settings=settings,
client=client
)
# Получаем метаинформацию (например, схему данных)
result_meta = get_query_result(query_id=query_id, result_index=0, client=client)
# Итерируемся по результату
result = read_query_result(query_id=query_id, result_index=0, client=client)
for row in result:
print(row)
Описание параметров конфигурации — в разделе Параметры конфигурации.
Через Spark Connect API
Основное отличие от классического Spark — только в способе создания сессии; остальной код остаётся прежним.
В примерах ниже показано, как создать Spark-сессию в каждом из режимов запуска:
import spyt
from yt.wrapper import YtClient
from spyt.connect import start_connect_server, wait_for_spark_connect_endpoint
from pyspark.sql import SparkSession
client = YtClient(proxy="<cluster-proxy>", token="<your-token>")
operation = start_connect_server(client)
endpoint = wait_for_spark_connect_endpoint(client, operation.id)
try:
spark = SparkSession.builder.remote(f"sc://{endpoint}").getOrCreate()
df = spark.read.format("yt").load("yt:///home/my-table")
df.show()
finally:
if spark:
spark.stop()
import spyt
from yt.wrapper import YtClient
from spyt.connect import start_connect_server_inner_cluster
from pyspark.sql import SparkSession
client = YtClient(proxy="<cluster-proxy>", token="<your-token>")
endpoint = start_connect_server_inner_cluster(client, discovery_path)
try:
spark = SparkSession.builder.remote(f"sc://{endpoint}").getOrCreate()
df = spark.read.format("yt").load("yt:///home/my-table")
df.show()
finally:
if spark:
spark.stop()
Параметры конфигурации QT
Параметры применяются только при работе через Query Tracker и передаются в формате JSON — в поле Settings (UI) или в словаре settings (API).
| Параметр | Описание | По умолчанию |
|---|---|---|
driver_cores |
CPU для драйвера | 1 |
driver_memory |
Память для драйвера | 1.5G |
num_executors |
Число экзекьюторов | 2 |
executor_cores |
CPU на один экзекьютор | 1 |
executor_memory |
Память на один экзекьютор | 4G |
spark_conf |
Конфигурационные параметры для запуска Spark-приложения. Можно указывать как стандартные параметры Spark, так и параметры SPYT | — |
Миграция с Livy
Начиная с SPYT 2.10.0 и Query Tracker 0.4 интеграция через Livy больше не поддерживается. SPYT Connect является её заменой.
Основное отличие от Livy: в Livy все драйверы запускались на одной машине с сервером, что ограничивало число одновременных сессий и не позволяло гибко настраивать ресурсы. В SPYT Connect каждый пользователь запускает Spark-драйвер как отдельную YTsaurus-операцию с индивидуальной конфигурацией.
| Характеристика | Livy | SPYT Connect |
|---|---|---|
| Одновременные сессии | Ограничены: все драйверы запускались на одном сервере, при нагрузке пользователи ждали в очереди | Не ограничены: каждый пользователь получает отдельную операцию YTsaurus |
| Изоляция сессий | Общий драйвер — нагрузка одного пользователя могла влиять на других | Каждый пользователь работает в изолированной среде выполнения |
| Квотирование | Отсутствует | Запросы выполняются в пользовательском пуле YTsaurus со стандартным квотированием |
| Конфигурация ресурсов | Фиксированная, задаётся на уровне сервера | Гибкая: CPU и память настраиваются отдельно для драйвера и каждого экзекьютора |
| Выбор пула ресурсов | Недоступен | Можно явно указать пул, в котором будут выполняться задачи |
Версии и совместимость
| Возможность | Минимальная версия |
|---|---|
| SPYT Connect с прямым сабмитом | SPYT 2.8.0 |
| SPYT Connect с внутренним кластером | SPYT 2.9.0 |
| Поддержка Spark 4.x | SPYT 2.10.0 |
| Замена Livy в Query Tracker | SPYT 2.10.0 + Query Tracker 0.4 |
Протокол Spark Connect появился в Spark 3.4, поэтому более ранние версии Spark не поддерживаются. Текущий SPYT Connect работает начиная с версии Spark 3.5.x.
Примечание
Для работы через Spark Connect API потребуются пакеты ytsaurus-spyt и pyspark-client.