Async Request в YTsaurus Flow (Java)
Данный пример пайплайна реализует событийно-ориентированный цикл запрос–ответ. Входящие события порождают запросы к обработчику, ответы возвращаются обратно в тот же компьютейшен и накапливаются во внешнем стейте. Пример демонстрирует циклическую топологию пайплайна и совместное использование внешнего стейта.
Компоненты
StateKeeperFunction
Обрабатывает два входных стрима — event и response. По событию event генерирует запрос с уникальным request_id и эмитирует его в стрим request. По событию response накапливает суммарную длину ответов в поле total_length внешнего стейта:
RequestProcessorFunction
Stateless-компьютейшен: получает запрос из стрима request, вычисляет длину строки и немедленно отправляет ответ в стрим response:
Регистрация компьютейшенов
Компьютейшены state и processor регистрируются аннотацией @FlowComputation на классах их process-функций:
PipelineMain
Единственная точка входа (запускает пайплайн или обслуживает его как компаньон — по YT_FLOW_MODE):
Ключевые паттерны
- Циклическая топология — стрим
responseвозвращается обратно вstate-компьютейшен, замыкая циклevent → request → response → state. Flow поддерживает такие графы явно. - Маршрутизация по
streamId— одна функция обрабатывает несколько входных стримов, определяя тип сообщения черезmessage.getStreamId(). - ExternalStateAccessor с PayloadBuilder — поле
total_lengthобновляется точечно:current.toBuilder()→ изменение →stateAccessor.set(updated.finish()). - Конфигурация через Spring Boot — компьютейшены регистрируются аннотацией
@FlowComputation;flow-spring-boot-starterуправляет жизненным циклом gRPC-сервера.