Shuffle в YTsaurus Flow (Java)

Пайплайн читает поток событий, группирует их по ключу и подсчитывает количество уникальных событий с использованием внешнего стейта (ExternalStateAccessor). Пример демонстрирует конфигурацию компаньона через Spring Boot.

Исходный код (Java)

Исходный код (Kotlin)

Компоненты компаньона

EventMapper

Процессная функция для source-компьютейшена reader. Выполняет парсинг и трансформацию входных данных. Ниже показана упрощённая версия — в реальном коде дополнительно выполняется JSON-парсинг поля data с помощью Jackson ObjectMapper:

EventReducer

Процессная функция с использованием ExternalStateAccessor для подсчета количества событий:

Логика работы:

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

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

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

PipelineMain

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

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

  • Конфигурация через Spring Boot — компьютейшены регистрируются аннотациями @FlowSourceComputation / @FlowComputation; flow-spring-boot-starter управляет жизненным циклом gRPC-сервера.
  • ExternalStateAccessor — работа с внешним стейтом через Payload и PayloadBuilder.
  • SourceComputation с ProcessFunction — reader использует EventMapper для трансформации входных данных на стороне компаньона.

См. также

Предыдущая
Следующая