Async Request в YTsaurus Flow (Java)

Данный пример пайплайна реализует событийно-ориентированный цикл запрос–ответ. Входящие события порождают запросы к обработчику, ответы возвращаются обратно в тот же компьютейшен и накапливаются во внешнем стейте. Пример демонстрирует циклическую топологию пайплайна и совместное использование внешнего стейта.

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

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

Компоненты

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-сервера.

См. также

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