Поддержанные конструкции 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)
- Джойн со статическими таблицами
- Джойн по префиксу ключевых колонок
- Джойн нескольких потоков между собой