---
metadata:
  - name: generator
    content: Diplodoc Platform v5.50.6
alternate:
  - https://ytsaurus.tech/docs/en/flow/cpp/examples/shuffle.md
  - https://ytsaurus.tech/docs/ru/flow/cpp/examples/shuffle.md
---
> **Documentation Index:** Fetch the complete configuration index at https://ytsaurus.tech/docs/ru/llms.txt

<!-- source: ru/_includes/flow/cpp/examples/shuffle.md -->
# Shuffle в YTsaurus Flow (C++)

[Пайплайн](https://github.com/ytsaurus/ytsaurus/tree/main/yt/yt/flow/examples/cpp/shuffle) читает поток в формате `JSON` из сортированной динамической таблицы, многократно группирует по разным ключам, а после считает число уникальных ключей во всех получившихся потоках. Для пайплайна также описан [тест](https://github.com/ytsaurus/ytsaurus/tree/main/yt/yt/flow/examples/cpp/shuffle/test/test_shuffle.py).

Пайплайн не пытается решить какую-либо бизнесовую задачу.

## Общее описание пайплайна

### Reader. Чтение входных данных

Первый компьютейшн - `Reader`. Он занимается чтением и первичным преобразованием данных из очереди. Вот части [спеки](https://ytsaurus.tech/docs/ru/flow/concepts/glossary.md#spec-and-dynamic-spec), связанные с данным компьютейшном:

```yson
{
    "spec" = {
        "computations" = {
            "reader" = {
                "computation_class_name" = "NYT::NFlow::TProcessFunctionSourceComputation";
                "processing_function" = "NYT::NFlow::NExample::TQueueReader";
                "output_stream_ids" = ["event"];
                "source_streams" = {
                    "queue" = {
                        "source_class_name" = "NYT::NFlow::TQueueSource";
                        "parameters" = {
                            "queue_path" = "<cluster=cluster_name>//path/to/queue";
                            "consumer_path" = "<cluster=cluster_name>//path/to/consumer";
                            "finite" = false;
                        };
                    };
                };
            };
        };
        "streams" = {
            "event" = {
                "schema" = [
                    {"name" = "key_a"; "type" = "uint64";};
                    {"name" = "key_b"; "type" = "uint64";};
                    {"name" = "key_c"; "type" = "uint64";};
                    {"name" = "key_d"; "type" = "uint64";};
                    {"name" = "value"; "type" = "string";};
                ];
            }
        };
    };
}
```

Разберем детально.

- `computations/reader/source_streams` содержит источник `queue` с типом `NYT::NFlow::TQueueSource`. Данный `source` предназначен для чтения данных из сортированной динамической таблицы с использованием `consumer`. В `parameters` указывается из какой очереди и каким консьюмером необходимо читать данные. С параметрами детальнее можно познакомиться в рамках класса `NYT::NFlow::TQueueSourceParameters`.
- Для управления `Computation` необходимо использовать `NYT::NFlow::TQueueSourceController` - так как нам нужно определять число и настройки [партиций](https://ytsaurus.tech/docs/ru/flow/concepts/glossary.md#partition) на базе входной сортированной динамической таблицы.
- `streams` содержит один поток `event` - распаршенный поток на выходе из `reader`, доступный другим `Computation`. Для него описана соответствующая схема. Этот же поток зарегистрирован и в `computations/reader/output_stream_ids`.
- Класс `TQueueReader` реализует `IProcessFunction`; его запускает `TProcessFunctionSourceComputation`, указанный в `computation_class_name`. Встроенный passthrough здесь не подходит, поскольку для разбора `JSON` нужна пользовательская реализация `ProcessMessage`.

{% code '/yt/yt/flow/examples/cpp/shuffle/lib/shuffle_functions.cpp' lang='cpp' lines='[BEGIN example_shuffle_queue_reader]-[END example_shuffle_queue_reader]' %}

- Source-адаптер работает в Swift-режиме: выходные потоки не материализуются в YTsaurus, сохраняется только метаинформация, необходимая для детерминированной работы.
- Доступом к очереди, включая нелокальный кластер, управляют `TQueueSource` и адаптер; process function получает уже прочитанное сообщение и не обращается к клиенту YTsaurus напрямую.

### Shuffle

В пайплайне присутствует несколько перемешиваний: `shuffle_a`, `shuffle_b`, `shuffle_c`, `shuffle_d`. Каждый из них группирует входной поток по соответствующему ключу - `key_a`, `key_b`, `key_c` или `key_d`. Никаких преобразований с данными они не делают, лишь демонстрируют возможность сгруппировать разные объекты.

Разберем спеку на примере`shuffle_b`:

```yson
{
    "spec" = {
        "computations" = {
            "shuffle_b" = {
                "computation_class_name" = "NYT::NFlow::TSwiftPassthroughComputation";
                "group_by_schema" = [
                    {"name" = "hash"; "expression" = "farm_hash(key_b)"; "type" = "uint64";};
                    {"name" = "key_b"; "type" = "uint64";};
                ];
                "input_stream_ids" = ["event_a"];
                "output_stream_ids" = ["event_b"];
            };
        };
        streams = {
            "event_a" = {
                "schema" = [
                    {"name" = "key_a"; "type" = "uint64";};
                    {"name" = "key_b"; "type" = "uint64";};
                    {"name" = "key_c"; "type" = "uint64";};
                    {"name" = "key_d"; "type" = "uint64";};
                    {"name" = "value"; "type" = "string";};
                ];
            };
            "event_b" = {
                "schema" = [
                    {"name" = "key_a"; "type" = "uint64";};
                    {"name" = "key_b"; "type" = "uint64";};
                    {"name" = "key_c"; "type" = "uint64";};
                    {"name" = "key_d"; "type" = "uint64";};
                    {"name" = "value"; "type" = "string";};
                ];
            };
        };
    };
}
```

- Так как в рамках примера нет какого-либо преобразования данных, то нам достаточно `NYT::NFlow::TSwiftPassthroughComputation`. Если преобразование понадобится, пользовательскую логику следует реализовать как `IProcessFunction` и запустить через `NYT::NFlow::TProcessFunctionSwiftMapComputation`.
- `NYT::NFlow::TSwiftPassthroughComputation` не материализует данные в YTsaurus.
- `group_by_schema` содержит соответствующий ключ `key_b`. В него добавлена колонка `hash`, так как партиционирование во `Flow` работает только в предположении, что первая колонка содержит равномерно распределенные значения типа `uint64`.
- `input_stream_ids` и `output_stream_ids` содержат соответственно `event_a` и `event_b`.
- В `spec/streams` также содержатся `event_a` и `event_b` с описанием схемы целиком.

### Reduce

Последний `Computation`. Он читает потоки `event_a`, `event_b`, `event_c`, `event_d` и подсчитывает число встреч каждого `value`. Фактически, исходный поток обрабатывается четыре раза.

```yson
{
    "spec" = {
        "computations" = {
            "reducer" = {
                "computation_class_name" = "NYT::NFlow::TProcessFunctionComputation";
                "processing_function" = "NYT::NFlow::NExample::TReducer";
                "group_by_schema" = [
                    {"name" = "hash"; "expression" = "farm_hash(value)"; "type" = "uint64";};
                    {"name" = "value"; "type" = "string";};
                ];
                "input_stream_ids" = ["event_a"; "event_b"; "event_c"; "event_d";];
                "output_stream_ids" = [];
                "external_state_managers" = {
                    "/state" = {
                        "external_state_manager_class_name" = "NYT::NFlow::TSimpleExternalStateManager";
                        "parameters" = {
                            "path" = "//path/to/state";
                        };
                    };
                };
            };
        };
    };
};
```

- Для описания логики используется process function `TReducer`, запущенная через `TProcessFunctionComputation`.
- Для работы со [стейтом](https://ytsaurus.tech/docs/ru/flow/concepts/glossary.md#state) используется `TSimpleExternalStateManager`, который предоставляет прямой доступ к таблице. `TReducer` хранит `TMutableStateKeyClient<TSimpleExternalState>` и привязывает его к `"/state"` в методе `Init(const IRuntimeInitContextPtr&)`.

{% code '/yt/yt/flow/examples/cpp/shuffle/lib/shuffle_functions.cpp' lang='cpp' lines='[BEGIN example_shuffle_reducer]-[END example_shuffle_reducer]' %}
<!-- endsource: ru/_includes/flow/cpp/examples/shuffle.md -->

<!-- source: ru/_includes/flow/cpp/examples/shuffle_also.md -->
## См. также

- [Быстрый старт (C++)](https://ytsaurus.tech/docs/ru/flow/cpp/getting-started.md)
- [Computation (C++)](https://ytsaurus.tech/docs/ru/flow/cpp/computation.md)

- Для группировки по `value` обязательно необходимо указать эту колонку (и хэш от неё) в `group_by_schema`.
- В `input_stream_ids` перечисляются все потоки: `event_a`, `event_b`, `event_c`, `event_d` &mdash; чтобы читать все получившиеся потоки. С точки зрения "бизнес логики" это не самое осмысленное действие, однако исходной целью данного пайплайна было протестировать гарантии `exactly-once` даже в случае `Swift` цепочки.
- `TReducer` реализует `IProcessFunction` и запускается через `TProcessFunctionComputation`. Адаптер сохраняет `input_messages` и `output_messages` в YTsaurus, но в этом примере выходных сообщений нет. По сути, пайплайн сохраняет метаинформацию `reader`, метаинформацию (`message_id` и `key`) каждого входного сообщения `reducer` и таблицу `value => count`. Промежуточные passthrough-компьютейшены с YTsaurus не взаимодействуют.

### DynamicSpec

- Поле `dynamic_spec/computations/<computation_id>/parameters/desired_partition_count` заполняется для каждого `computation`, кроме `reader`. В рамках теста `test_shuffle.py` происходит изменение числа партиций.
- В `dynamic_spec/job_tracker/job_threads` указывается необходимое число тредов для выполнения всех джобов.

### Config для запуска

- Ключевое для запуска: `cluster_url`, `proxy_role`, `path`, `rpc_port`, `monitoring_port`.
- `controller/scheduler_period` выставлен в 200 для конкретного теста - в реальности должно быть достаточно дефолтного значения.
- `logging` - настройки логирования.

```yson
{
    "cluster_url" = "cluster_name";
    "path" = "//path/to/pipeline";
    "rpc_port" = 81;
    "monitoring_port" = 80;
    "controller" = {
        "scheduler_period" = 200;
    };
    "logging" = {
        "suppressed_messages" = [
        ];
        "rules" = [
            {
                "exclude_categories" = [
                    "Bus";
                    "Dns";
                    "Concurrency";
                    "QueryClient";
                    "Profiling";
                    "RpcClient";
                    "Monitoring";
                    "Net";
                    "Solomon";
                    "Jaeger";
                    "RpcProxyClient";
                    "RpcServer";
                    "Dns";
                    "BufferMetrics";
                ];
                "min_level" = "debug";
                "writers" = [
                    "Stderr";
                ];
            };
        ];
        "writers" = {
            "Stderr" = {
                "type" = "file";
                "file_name" = "/path/to/file.log";
            };
        };
    }
}
```
<!-- endsource: ru/_includes/flow/cpp/examples/shuffle_also.md -->
