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):

  1. Сохраняет запрос в стейт с FailedAttempts = 0
  2. Вызывает 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.

Структура пайплайна

  1. events → TStateKeeper → request (генерация запроса)
  2. request → TRequestProcessor → response (обработка с ретраями)
  3. response → TStateKeeper → стейт (накопление результатов)

В спеке для TRequestProcessor необходимо зарегистрировать секцию timer_streams для поддержки повторных попыток. Класс process function указывается в поле processing_function.

Исходный код

TRequestProcessor

TStateKeeper

См. также

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