Async Request в YTsaurus Flow (Go)
Пример пайплайна, реализующего асинхронный поход во внешний сервис: один компьютейшен превращает события в запросы и накапливает ответы во внешнем стейте, другой обслуживает запросы без стейта. Go-реализация того же сценария, что и C++ пример.
Структура
Компаньон обслуживает два компьютейшена, injector остаётся нативным сорсом из спеки:
-
state(stateKeeper) — stateful-компьютейшен, сгруппированный поkey, который:- принимает события из стрима
eventи порождает запрос в стримrequestсо случайнымrequest_id; - принимает ответы из стрима
responseи складывает суммарную длину (total_length) во внешний стейт/state.
- принимает события из стрима
-
processor(requestProcessor) — stateless-компьютейшен, сгруппированный поrequest_id: принимает запросы из стримаrequestи сразу отвечает длиной строки запроса в стримresponse.
Цикл event → request → response → state замыкается между двумя компьютейшенами. Событие отвечается запросом, а не сразу результатом, поэтому обслуживающая сторона никогда не задерживает обработку: ответ приходит позже, отдельным сообщением, и только тогда стейт ключа сдвигается.
main.go
Точка входа: создание пайплайна и регистрация обоих компьютейшенов.
state_keeper.go
Маршрутизация входных стримов (event / response) и работа с внешним стейтом.
request_processor.go
Stateless-обработчик запросов: вычисляет длину строки запроса и возвращает ответ.
Ключевые паттерны
- Маршрутизация по
msg.StreamID:switchпо идентификатору входного стрима позволяет одному компьютейшену обрабатывать несколько входов с разной логикой. Неизвестный стрим — ошибка, а не молчаливое игнорирование. - Случайный
request_id:rand.Uint64()связывает запрос с ответом. Запрос несёт ключ исходного события, поэтому ответ, партиционированный поrequest_id, возвращается к тому стейту, которому принадлежит поход. - Внешний стейт через
flow.OpenExternalState(rt, "/state", msg): строка преобразуется вtotalLengthStateчерезConvertTo, изменяется как структура и сохраняется черезConvertFrom. - Stateless-компьютейшен:
requestProcessorне использует стейт и сгруппирован поrequest_id, а не по ключу события, поэтому запросы одного ключа расходятся по всем партициям и масштабируются независимо. - Зависимость стримов:
streams_dependencyв спеке объявляет, чтоrequestпорождается изevent— воркер учитывает это при продвижении вотермарков.