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

See also