Retryable Async Request in YTsaurus Flow (Java)
The pipeline extends AsyncRequest: the request handler now supports retry attempts with a delay. If a failure occurs, RequestProcessorFunction saves the request state in the internal state and sets a timer for 5 seconds; success is determined by the predicate (requestId + failedAttempts) % 3 == 0.
Components
RequestProcessorFunction
This implements the retry logic using an internal YSON state and timers. When you receive a request, you save it to the state and call tryRequest. When the timer fires, you load the state and retry the attempt:
The helper method tryRequest checks the success predicate. If the attempt fails, it increments the attempt counter, saves the state, and schedules the next attempt.
RequestState
This is a YSON-serializable model of the request state. It stores requestId, key, the request text request, and the counter failedAttempts:
StateKeeperFunction
This is identical to the StateKeeperFunction from AsyncRequest: it processes the event and response streams and accumulates total_length in the external state. The retry logic is fully encapsulated in RequestProcessorFunction.
Registering computations
You register the state and processor computations with the @FlowComputation annotation on the classes of their process functions:
PipelineMain
This single entry point starts the pipeline or serves as its companion, depending on YT_FLOW_MODE.
Key patterns
- Retries via timers: if a failure occurs,
output.addTimer(now + 5, 0L)schedules a repeated call toonTimer; the state between attempts is stored inYsonStateAccessor. - Success predicate:
(requestId + failedAttempts) % 3 == 0simulates an unstable external service; in real tasks, you replace this with a check of the HTTP status or another indicator. - Clearing the state after success:
accessor.clear()is called only on a successful response, preventing reprocessing. - Separation of concerns: the session state logic (
StateKeeperFunction) is separated from the retry logic (RequestProcessorFunction), which makes it easier to test each part independently.
Differences from AsyncRequest
| Aspect | AsyncRequest | RetryableAsyncRequest |
|---|---|---|
RequestProcessorFunction |
Stateless, responds immediately | Has internal state and timers |
| Failure handling | Not provided | Retry with a 5-second delay |
| Request state | Absent | RequestState in YSON |