Supported YQL constructs in YTsaurus Flow
Note
This page describes YQL constructs specific to working with Flow. The general YQL syntax is described in the YQL documentation.
Row-wise stream transformation (map)
The main construct is INSERT INTO ... SELECT ... FROM. It reads a stream from an ordered dynamic table (a YTsaurus queue), applies a transformation, and writes the result to another queue.
Example of a query part:
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;
You can work with files and UDFs (built-in and user-defined), lambda expressions, and code generation.
You can combine multiple INSERT INTO ... SELECT operations in a single query.
Stream join with a dynamic table (lookup join)
You can join a stream with a sorted dynamic table (a key-value table). Supported join types for the “stream + KV table” pair are LEFT, LEFT ONLY, LEFT SEMI, and INNER. For the “KV table + stream” pair, the types are symmetric.
$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;
Roadmap
Under development:
- Aggregations over fixed-size windows (hopping windows)
- Join with static tables
- Join by prefix of key columns
- Join of multiple streams with each other