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

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

[Пайплайн](https://ytsaurus.tech/docs/ru/flow/concepts/glossary.md#pipeline) читает [поток](https://ytsaurus.tech/docs/ru/flow/concepts/glossary.md#stream-and-computation) событий, группирует их по ключу и подсчитывает количество уникальных событий с использованием внешнего [стейта](https://ytsaurus.tech/docs/ru/flow/concepts/glossary.md#state) (ExternalStateAccessor). Пример демонстрирует конфигурацию [компаньона](https://ytsaurus.tech/docs/ru/flow/concepts/glossary.md#companion) через Spring Boot.

[Исходный код (Java)](https://github.com/ytsaurus/ytsaurus/tree/main/yt/yt/flow/examples/java/shuffle)

[Исходный код (Kotlin)](https://github.com/ytsaurus/ytsaurus/tree/main/yt/yt/flow/examples/kotlin/shuffle)
## Компоненты компаньона

### EventMapper

Процессная функция для source-[компьютейшена](https://ytsaurus.tech/docs/ru/flow/concepts/glossary.md#stream-and-computation) `reader`. Выполняет парсинг и трансформацию входных данных. Ниже показана упрощённая версия — в [реальном коде](https://github.com/ytsaurus/ytsaurus/tree/main/yt/yt/flow/examples/java/shuffle) дополнительно выполняется JSON-парсинг поля `data` с помощью Jackson `ObjectMapper`:

{% list tabs group=lang %}

- Java

  {% code '/yt/yt/flow/examples/java/shuffle/shuffle/src/main/java/tech/ytsaurus/flow/examples/shuffle/EventMapper.java' lang='java' lines='[BEGIN on_message]-[END on_message]' keep-indents %}

- Kotlin

  {% code '/yt/yt/flow/examples/kotlin/shuffle/shuffle/src/main/kotlin/tech/ytsaurus/flow/examples/shuffle/EventMapper.kt' lang='kotlin' lines='[BEGIN on_message]-[END on_message]' keep-indents %}

{% endlist %}

### EventReducer

Процессная функция с использованием [ExternalStateAccessor](https://ytsaurus.tech/docs/ru/flow/java/state.md#external-state) для подсчета количества событий:

{% list tabs group=lang %}

- Java

  {% code '/yt/yt/flow/examples/java/shuffle/shuffle/src/main/java/tech/ytsaurus/flow/examples/shuffle/EventReducer.java' lang='java' lines='[BEGIN on_message]-[END on_message]' keep-indents %}

- Kotlin

  {% code '/yt/yt/flow/examples/kotlin/shuffle/shuffle/src/main/kotlin/tech/ytsaurus/flow/examples/shuffle/EventReducer.kt' lang='kotlin' lines='[BEGIN on_message]-[END on_message]' keep-indents %}

{% endlist %}

Логика работы:
1. Получаем `ExternalStateAccessor` для стейта `"shuffle-state"`, привязанного к ключу текущего сообщения.
2. Извлекаем текущее значение стейта. Если стейта нет — `getOrDefault()` вернет пустой `Payload`.
3. Создаем `PayloadBuilder` из текущего стейта, увеличиваем счетчик.
4. Сохраняем обновленный стейт.

### Регистрация компьютейшенов

Компьютейшен-источник `reader` регистрируется аннотацией `@FlowSourceComputation`, а трансформация `reducer` — аннотацией `@FlowComputation`:

{% list tabs group=lang %}

- Java

  {% code '/yt/yt/flow/examples/java/shuffle/shuffle/src/main/java/tech/ytsaurus/flow/examples/shuffle/EventMapper.java' lang='java' lines='[BEGIN registration]-[END registration]' %}

  {% code '/yt/yt/flow/examples/java/shuffle/shuffle/src/main/java/tech/ytsaurus/flow/examples/shuffle/EventReducer.java' lang='java' lines='[BEGIN registration]-[END registration]' %}

- Kotlin

  {% code '/yt/yt/flow/examples/kotlin/shuffle/shuffle/src/main/kotlin/tech/ytsaurus/flow/examples/shuffle/EventMapper.kt' lang='kotlin' lines='[BEGIN registration]-[END registration]' %}

  {% code '/yt/yt/flow/examples/kotlin/shuffle/shuffle/src/main/kotlin/tech/ytsaurus/flow/examples/shuffle/EventReducer.kt' lang='kotlin' lines='[BEGIN registration]-[END registration]' %}

{% endlist %}

### PipelineMain

Единственная точка входа (запускает пайплайн или обслуживает его как компаньон — по `YT_FLOW_MODE`):

{% list tabs group=lang %}

- Java

  {% code '/yt/yt/flow/examples/java/shuffle/shuffle/src/main/java/tech/ytsaurus/flow/examples/shuffle/PipelineMain.java' lang='java' lines='[BEGIN main]-[END main]' keep-indents %}

- Kotlin

  {% code '/yt/yt/flow/examples/kotlin/shuffle/shuffle/src/main/kotlin/tech/ytsaurus/flow/examples/shuffle/PipelineMain.kt' lang='kotlin' lines='[BEGIN main]-[END main]' keep-indents %}

{% endlist %}

## Ключевые паттерны

- **Конфигурация через Spring Boot** — компьютейшены регистрируются аннотациями `@FlowSourceComputation` / `@FlowComputation`; `flow-spring-boot-starter` управляет жизненным циклом gRPC-сервера.
- **ExternalStateAccessor** — работа с внешним стейтом через `Payload` и `PayloadBuilder`.
- **SourceComputation с ProcessFunction** — `reader` использует `EventMapper` для трансформации входных данных на стороне компаньона.
<!-- endsource: ru/_includes/flow/java/examples/shuffle.md -->

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

- [Быстрый старт (Java)](https://ytsaurus.tech/docs/ru/flow/java/getting-started.md)
- [Computation (Java)](https://ytsaurus.tech/docs/ru/flow/java/computation.md)
<!-- endsource: ru/_includes/flow/java/examples/shuffle_also.md -->
