Quick start with YTsaurus Flow (YQL)

Use YQL over Flow to describe a pipeline for streaming data processing as a declarative SQL query — without writing code in C++, Java, Python, or Go. The pipeline runs as a vanilla operation on the selected YTsaurus cluster.

Warning

This feature is under active development, and not all planned functionality is available yet.

Pragmas

You control a YQL over Flow query with a set of pragmas:

Pragma Description
PRAGMA Engine = "ytflow"; Selects the Flow engine to run the query
PRAGMA Ytflow.Cluster = "..."; Cluster for the pipeline’s internal tables and output ordered queues
PRAGMA Ytflow.RuntimeCluster = "..."; Cluster to run the vanilla operation.
PRAGMA Ytflow.PipelineDirectory = "..."; Path to the directory with pipelines in YTsaurus
PRAGMA Ytflow.PipelineName = "..."; Pipeline name. Full path: {pipeline_directory}/{pipeline_name}
PRAGMA Ytflow.WorkerCount = "..."; Number of worker jobs for the vanilla operation
PRAGMA Ytflow.EnableComputationPatternResources = "true"; Reuses computation patterns across graphs in one worker. Default is false

First query

Example: row-by-row transformation of a stream (map).

-- select the Flow engine
PRAGMA Engine = "ytflow";

-- cluster for the pipeline’s internal tables
PRAGMA Ytflow.Cluster = "<cluster-name>";
-- cluster for the vanilla operation
PRAGMA Ytflow.RuntimeCluster = "<cluster-name>";
-- directory with pipelines
PRAGMA Ytflow.PipelineDirectory = "//home/my-project/pipelines";
-- pipeline name
PRAGMA Ytflow.PipelineName = "my-pipeline";
-- number of workers
PRAGMA Ytflow.WorkerCount = "1";

-- read from the input queue, transform, write to the output queue
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;

The query starts a pipeline that continuously processes messages from the input queue and writes the results to the output queue. The output table schemas are automatically inferred from the query.

For a description of all supported YQL constructs, see the Supported constructs section.

How to run

Prerequisites

You need read and write permissions for all directories mentioned in the query, as well as a compute quota on the YTsaurus cluster specified as Ytflow.RuntimeCluster.

You can run the query in two ways:

Via the YTsaurus UI: open the Queries tab on the runtime cluster and run the query.

Via the Python client:

from yt.wrapper import YtClient

# any production cluster
client = YtClient('<cluster-name>')

# run the query and wait for completion
client.run_query(
    engine='yql',
    settings=dict(
        # pass the runtime cluster here
        cluster='<cluster-name>',
    ),
    query='<YQL query>',
    sync=True,
)

After the query finishes, a pipeline starts on the cluster and runs continuously. If a pipeline with the same name already exists, it stops and finishes processing all internal streams, then the new version starts.

Monitoring

To track the running pipeline, you can use:

  • Dashboard — Flow → Monitoring tab.
  • Controller logs (worker status, possible issues):
    yt --proxy <pipeline-cluster> flow show-logs //home/my-project/pipelines/my-pipeline
    
  • Job logs — via the vanilla operation, which is available through the link from the flowPublish cube in the pipeline graph.

See also