Быстрый старт в YTsaurus Flow (YQL)

YQL over Flow позволяет описать пайплайн потоковой обработки данных в виде декларативного SQL-запроса — без написания кода на C++, Java, Go или Python. Пайплайн запускается как ванилла операция на выбранном кластере YTsaurus.

Важно

Находится в активной разработке, ещё не вся запланированная функциональность доступна.

Полезные ссылки

Прагмы

Запрос YQL over Flow управляется через набор прагм:

Прагма Описание
PRAGMA Engine = "ytflow"; Выбор движка Flow для выполнения запроса
PRAGMA Ytflow.Cluster = "..."; Кластер для внутренних таблиц пайплайна и выходных упорядоченных очередей
PRAGMA Ytflow.RuntimeCluster = "..."; Кластер для запуска ванилла операции.
PRAGMA Ytflow.PipelineDirectory = "..."; Путь к директории с пайплайнами в YTsaurus
PRAGMA Ytflow.PipelineName = "..."; Имя пайплайна. Полный путь: {pipeline_directory}/{pipeline_name}
PRAGMA Ytflow.WorkerCount = "..."; Количество воркер-джобов ванилла операции
PRAGMA Ytflow.EnableComputationPatternResources = "true"; Включает переиспользование шаблонов вычислений между графами одного воркера. По умолчанию false

Первый запрос

Пример: построчное преобразование стрима (мап).

-- выбрать движок Flow
PRAGMA Engine = "ytflow";

-- кластер для внутренних таблиц пайплайна
PRAGMA Ytflow.Cluster = "<cluster-name>";
-- кластер для ванилла операции
PRAGMA Ytflow.RuntimeCluster = "<cluster-name>";
-- директория с пайплайнами
PRAGMA Ytflow.PipelineDirectory = "//home/my-project/pipelines";
-- имя пайплайна
PRAGMA Ytflow.PipelineName = "my-pipeline";
-- число воркеров
PRAGMA Ytflow.WorkerCount = "1";

-- читать из входной очереди, трансформировать, писать в выходную
INSERT INTO
    <cluster-name>.`//home/my-project/output_queues/sink_queue`
SELECT
    string_field || "_processed" AS string_field,
    int64_field,
    EndsWith(string_field, "bar") AS predicate
FROM
    <cluster-name>.`//home/my-project/input_queues/source_queue`
WHERE int64_field > 1;

Запрос запускает пайплайн, который непрерывно обрабатывает сообщения из входной очереди и пишет результаты в выходную. Схемы выходных таблиц выводятся автоматически из запроса.

Описание всех поддержанных конструкций YQL см. в разделе Поддержанные конструкции.

Как запустить

Пререквизиты

Нужно иметь права на чтение и запись во все упоминаемые в запросе директории, а также вычислительную квоту на кластере YTsaurus, указанном как Ytflow.RuntimeCluster.

Есть два способа запустить запрос:

Через UI YTsaurus: откройте вкладку Queries на рантайм кластере и выполните запрос.

Через Python-клиент:

from yt.wrapper import YtClient

# любой продакшн кластер
client = YtClient('<cluster-name>')

# запустить запрос и дождаться завершения
client.run_query(
    engine='yql',
    settings=dict(
        # рантайм кластер передаётся здесь
        cluster='<cluster-name>',
    ),
    query='<YQL query>',
    sync=True,
)

После завершения запроса на кластере запустится пайплайн, который будет выполняться непрерывно. Если пайплайн с таким именем уже существует — он остановится с дообработкой всех внутренних потоков, после чего запустится новая версия.

Мониторинг

Для отслеживания работы запущенного пайплайна доступны:

  • Дашборд — вкладка Flow → Monitoring.
  • Логи контроллера (состояние воркеров, возможные проблемы):
    yt --proxy <кластер-пайплайна> flow show-logs //home/my-project/pipelines/my-pipeline
    
  • Логи джобов — через ванилла операцию, доступную по ссылке из кубика flowPublish в графе пайплайна.

См. также