Async Request in YTsaurus Flow (Python)

This example shows a two-component pipeline that implements asynchronous request processing. One computation routes events into requests and accumulates responses. The other performs computation without state. This is a Python implementation of a similar C++ example.

Source code

Structure

The pipeline includes two computations:

  1. state (StateKeeperFunction) — a stateful computation that:

    • Accepts events from the event stream and generates requests to the request stream with a unique request_id.
    • Accepts responses from the response stream and accumulates the total length (total_length) in the external state.
  2. processor (RequestProcessorFunction) — a stateless computation that accepts requests from the request stream and immediately returns a response (the length of the request string) to the response stream.

The event → request → response → state cycle closes between the two computations.

state_keeper_function.py

This file handles routing of incoming streams (event / response) and working with the external state.

request_processor_function.py

This is a stateless request handler: it calculates the length of the request string and returns the response.

__main__.py

This is the entry point: it creates the pipeline and registers both computations.

{% code '/yt/yt/flow/examples/python/async_request/main.py' lang='python' lines='[BEGIN main]-[END main]' %}

Key patterns

  • Routing by stream_id: the if stream_id == "event" / "response" branching lets a single computation handle multiple input streams with different logic.
  • Generating a unique request_id: random.getrandbits(64) ensures correlation between the request and the response in the asynchronous cycle.
  • External state via ctx.external_state("/state", message): the to_builder() / set() pattern supports cumulative updates to the external state.
  • Stateless computation: RequestProcessorFunction doesn’t use state — it’s a pure transformation of a request into a response, which lets you scale it independently.
  • Two-component pipeline: pipeline.add("state", ...) and pipeline.add("processor", ...) register the computations. The streams between them are described in the spec.

See also