Shuffle в YTsaurus Flow (Java)
Пайплайн читает поток событий, группирует их по ключу и подсчитывает количество уникальных событий с использованием внешнего стейта (ExternalStateAccessor). Пример демонстрирует конфигурацию компаньона через Spring Boot.
Компоненты компаньона
EventMapper
Процессная функция для source-компьютейшена reader. Выполняет парсинг и трансформацию входных данных. Ниже показана упрощённая версия — в реальном коде дополнительно выполняется JSON-парсинг поля data с помощью Jackson ObjectMapper:
EventReducer
Процессная функция с использованием ExternalStateAccessor для подсчета количества событий:
Логика работы:
- Получаем
ExternalStateAccessorдля стейта"shuffle-state", привязанного к ключу текущего сообщения. - Извлекаем текущее значение стейта. Если стейта нет —
getOrDefault()вернет пустойPayload. - Создаем
PayloadBuilderиз текущего стейта, увеличиваем счетчик. - Сохраняем обновленный стейт.
Регистрация компьютейшенов
Компьютейшен-источник reader регистрируется аннотацией @FlowSourceComputation, а трансформация reducer — аннотацией @FlowComputation:
PipelineMain
Единственная точка входа (запускает пайплайн или обслуживает его как компаньон — по YT_FLOW_MODE):
Ключевые паттерны
- Конфигурация через Spring Boot — компьютейшены регистрируются аннотациями
@FlowSourceComputation/@FlowComputation;flow-spring-boot-starterуправляет жизненным циклом gRPC-сервера. - ExternalStateAccessor — работа с внешним стейтом через
PayloadиPayloadBuilder. - SourceComputation с ProcessFunction —
readerиспользуетEventMapperдля трансформации входных данных на стороне компаньона.