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.
Structure
The pipeline includes two computations:
-
state(StateKeeperFunction) — a stateful computation that:- Accepts events from the
eventstream and generates requests to therequeststream with a uniquerequest_id. - Accepts responses from the
responsestream and accumulates the total length (total_length) in the external state.
- Accepts events from the
-
processor(RequestProcessorFunction) — a stateless computation that accepts requests from therequeststream and immediately returns a response (the length of the request string) to theresponsestream.
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: theif 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): theto_builder()/set()pattern supports cumulative updates to the external state. - Stateless computation:
RequestProcessorFunctiondoesn’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", ...)andpipeline.add("processor", ...)register the computations. The streams between them are described in the spec.