Async Request в YTsaurus Flow (C++)
Пайплайн демонстрирует паттерн асинхронных внешних запросов с использованием process functions. События поступают на вход, преобразуются в запросы, обрабатываются детерминированной функцией и результаты накапливаются в стейте.
Компоненты пайплайна
TStateKeeper
TStateKeeper реализует IProcessFunction и запускается адаптером TProcessFunctionComputation. Функция использует TSimpleExternalStateManager для работы с внешним стейтом и обрабатывает два вида входных сообщений:
- Поток
event: при получении события создаетTRequestMessageс уникальнымRequestIdи отправляет его в потокrequest. - Поток
response: при получении ответа обновляет стейт — суммируетtotal_lengthиз всех полученных ответов.
Различение потоков происходит через message->StreamId.
TRequestProcessor
TRequestProcessor реализует IProcessFunction и запускается адаптером TProcessFunctionSwiftMapComputation, который не сохраняет входные и выходные сообщения в YTsaurus. Функция получает TRequestMessage, выполняет обработку (в данном примере — вычисляет длину запроса) и генерирует TResponseMessage.
Использование SwiftMap-адаптера обосновано тем, что обработка запроса является чистой функцией: при одинаковых входных данных всегда генерируется одинаковый результат.
Типы сообщений
- TEventMessage — наследник
TYsonMessage. Содержит поляKeyиData. - TRequestMessage — наследник
TYsonMessage. Содержит поляRequestId,KeyиRequest. - TResponseMessage — наследник
TYsonMessage. Содержит поляRequestId,KeyиLength.
Все типы сообщений регистрируются через макрос YT_FLOW_DEFINE_YSON_MESSAGE.
Ключевой паттерн: цикл запрос-ответ
Основная идея данного примера — построение цикла запрос-ответ внутри пайплайна с использованием нескольких потоков:
- events →
TStateKeeper→ request (генерация запроса) - request →
TRequestProcessor→ response (обработка запроса) - response →
TStateKeeper→ стейт (накопление результатов)
TStateKeeper одновременно является и потребителем событий, и потребителем ответов. Он использует input_stream_ids = ["event", "response"] и определяет тип входного сообщения по StreamId.
В спеке process functions связываются с адаптерами явно:
- для
stateуказаныcomputation_class_name = "NYT::NFlow::TProcessFunctionComputation"иprocessing_function = "NYT::NFlow::NExample::TStateKeeper"; - для
processorуказаныcomputation_class_name = "NYT::NFlow::TProcessFunctionSwiftMapComputation"иprocessing_function = "NYT::NFlow::NExample::TRequestProcessor".
Управление стейтом
TStateKeeper использует TSimpleExternalStateManager для хранения суммы длин всех обработанных запросов. Клиент стейта (TMutableStateKeyClient<TSimpleExternalState>) привязывается в Init(const IRuntimeInitContextPtr&) через initContext->InitExternalStateClient(StateClient_, "/state"). Параметры стейта (path к таблице и т.п.) объявляются в секции external_state_managers спеки компьютейшена.
Функция main
В main регистрируются три потока, а process functions регистрируются через YT_FLOW_DEFINE_PROCESS_FUNCTION:
RegisterStream<TEventMessage>("event")— входные событияRegisterStream<TRequestMessage>("request")— запросы к процессоруRegisterStream<TResponseMessage>("response")— ответы от процессора