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 подключается к нему

Способы подключения

Способ

Когда подходит

Через UI Query Tracker

Подходит аналитикам и всем, кто работает с данными через интерфейс YTsaurus

Через Query Tracker API

Для автоматизации SQL-запросов из Python

Через Spark Connect API

Для тех, кто хочет использовать DataFrame API или управлять жизненным циклом драйвера вручную

Через UI Query Tracker

Чтобы выполнить SQL-запрос с использованием SPYT Connect:

  1. Откройте вкладку Queries в интерфейсе YTsaurus.
  2. В списке движков выберите SPYT.
  3. Введите SQL-запрос.
  4. В поле Settings укажите конфигурацию в формате JSON.
  5. Нажмите Run и дождитесь результата.

SPYT Connect в Query Tracker

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

{
  "discovery_path": "//home/spark/my-cluster"
}

SPYT Connect с внутренним кластером

Через 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.

Что дальше

Предыдущая
Следующая