Retryable Async Request в YTsaurus Flow (C++)
Пайплайн аналогичен Async Request, но добавляет логику повторных попыток (ретраев) с использованием таймеров и внутреннего YsonState.
Отличие от Async Request
Ключевое отличие — TRequestProcessor запускается адаптером TProcessFunctionComputation (а не TProcessFunctionSwiftMapComputation), поскольку ему необходимо:
- хранить внутренний стейт для отслеживания числа неудачных попыток;
- использовать таймеры для повторных попыток через заданный интервал.
Компоненты пайплайна
TRequestProcessor
TRequestProcessor реализует IProcessFunction и использует TMutableStateKeyClient<TDelayedRequestState> для хранения внутреннего стейта (Internal YsonState). Стейт инициализируется в Init(const IRuntimeInitContextPtr& initContext) через initContext->InitClient(RequestStateClient_, "request_state").
При обработке входного сообщения (ProcessMessage):
- Сохраняет запрос в стейт с
FailedAttempts = 0 - Вызывает
TryRequestдля выполнения попытки
Метод TryRequest содержит основную логику ретраев:
- Если запрос "не удался" (определяется через
IsRequestSuccessful), увеличивает счетчикFailedAttemptsи ставит таймер черезoutput->AddTimer(GetNextAttempt(context)) - Если запрос "удался", создает
TResponseMessage, сбрасывает стейт черезstate.Clear()и отправляет ответ
При срабатывании таймера (ProcessTimer) вызывается повторная попытка TryRequest с текущим стейтом.
TStateKeeper
Полностью аналогичен TStateKeeper из примера Async Request: принимает входные события и ответы, хранит аккумулированный результат во внешнем стейте.
Паттерн ретраев
Логика ретраев основана на следующих элементах:
- TDelayedRequestState — наследник
NYTree::TYsonStruct, хранитFailedAttemptsи самRequest - TMutableStateKeyClient — клиент для работы с внутренним YsonState. В отличие от
TSimpleExternalStateManager, стейт хранится во внутренних таблицах Flow, а не во внешней пользовательской таблице - Таймеры — при неудачной попытке устанавливается таймер с задержкой
Delayчерезoutput->AddTimer(GetNextAttempt(context)) - Константа
MaxRetries = 3задаёт период симуляции успешной попытки: запрос выполняется сразу или не более чем после двух ретраев.Delay = 5задаёт задержку между попытками
Вычисление времени следующей попытки выполняется через context->GetCurrentTimestamp(), что обеспечивает корректную работу с системным временем Flow.
Структура пайплайна
- events →
TStateKeeper→ request (генерация запроса) - request →
TRequestProcessor→ response (обработка с ретраями) - response →
TStateKeeper→ стейт (накопление результатов)
В спеке для TRequestProcessor необходимо зарегистрировать секцию timer_streams для поддержки повторных попыток. Класс process function указывается в поле processing_function.