Поддержанные конструкции YQL в YTsaurus Flow

Примечание

Данная страница описывает конструкции YQL, специфичные для работы с Flow. Общий синтаксис YQL описан в документации YQL.

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

Основная конструкция — INSERT INTO ... SELECT ... FROM. Читает стрим из упорядоченной динамической таблицы (очереди YTsaurus), применяет преобразование и записывает результат в другую очередь.

Пример части запроса:

INSERT INTO
    <cluster-name>.`//home/my-project/output_queues/sink_queue`
SELECT
    string_field || "_processed" AS processed_field,
    int64_field * 2 AS doubled,
    EndsWith(string_field, "bar") AS predicate
FROM
    <cluster-name>.`//home/my-project/input_queues/source_queue`
WHERE int64_field > 0;

Поддерживается работа с файлами и UDF (встроенными и пользовательскими), лямбда-выражениями и кодогенерацией.

В одном запросе можно комбинировать несколько операций INSERT INTO ... SELECT.

Джойн потока с динамической таблицей (lookup join)

Стрим можно джойнить с сортированной динамической таблицей (key-value таблицей). Поддерживаемые типы джойна для пары «поток + KV таблица»: LEFT, LEFT ONLY, LEFT SEMI, INNER. Для пары «KV таблица + поток» типы симметричны.

$input_stream =
    SELECT key, value, key || "_before" AS key_before
    FROM <cluster-name>.`//home/my-project/input_queues/source_queue`
    WHERE value > 2;

$joined_stream =
    SELECT
        left_arg.key AS key,
        left_arg.value AS value,
        left_arg.key_before,
        right_arg.kv_value
    FROM $input_stream AS left_arg
    INNER JOIN
        <cluster-name>.`//home/my-project/states/kv_table` AS right_arg
    USING (key);

INSERT INTO <cluster-name>.`//home/my-project/output_queues/sink_queue`
SELECT * from $joined_stream
WHERE value * 2 <= kv_value;

Ближайшие планы

В разработке:

  • Агрегации по окнам фиксированного размера (hopping windows)
  • Джойн со статическими таблицами
  • Джойн по префиксу ключевых колонок
  • Джойн нескольких потоков между собой

См. также

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