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.

Source code (Java)

Source code (Kotlin)

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 to onTimer; the state between attempts is stored in YsonStateAccessor.
  • Success predicate: (requestId + failedAttempts) % 3 == 0 simulates 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

See also