Stateful processing in YTsaurus Flow
Use stateful processing to handle events with read-modify-write operations on states stored in YTsaurus. For example, you can count statistics for incoming events by key: you load the old value, update it, and write it back.
State access model
You access state inside a computation through three components:
- State — user data tied to a key and stored in a YTsaurus dynamic table.
- State client — a type-specific object (
Client<TState>) that the process function creates for each named state and initializes inInit. You access state by key through the client. It can be read-write (works with both internal and external states) or read-only (joiner for external state). - State accessor — what the client returns for a specific key (based on an input message, timer, or explicit key): a representation of the state of the same type
TState. The accessor behaves like a smart pointer to the state: a read-write accessor lets you read, modify, and clear the state; a read-only accessor lets you only read it.
Warning
The accessor is valid only within the current epoch. Don't store it in process-function fields or reuse it across epochs — get the state through the client again in each epoch.
State types
Internal State
This is the simplest way to work with state: you don't need to create tables — Flow manages them automatically. Data loads at the start of the epoch and writes on commit. You get read-write access through the same client you use for external state. The state type can be arbitrary; the only requirement is that it's serializable to YSON (to persist between epochs).
External State
This is state in a user-managed dynamic table. You create and manage the tables. You get read-write access through the same client you use for internal state, but the backend is a state manager (TSimpleExternalStateManager). You declare it in the computation spec in the top-level external_state_managers section; the implementation resolves via external_state_manager_class_name. It supports caching.
External State Joiner
You get read-only access to external states via a key-based join — through a joiner (TSimpleExternalStateJoiner). You declare it in the computation spec in the top-level external_state_joiners section (at the same level as external_state_managers); the implementation resolves via external_state_joiner_class_name. It supports TTL-based caching: loaded states live in the shared StateCache and reload from YT only after the TTL expires or the state is evicted from the cache.
Warning
One table, one writer
Only one computation should write to an external state table: writes from different partitions and transactions break state consistency. The state manager owns its table for writing: TSimpleExternalStateManager declares it as its own. To get read-only access to another computation's state, use an external state joiner (TSimpleExternalStateJoiner) or send messages to the writer computation — joiners don't lock the table. The spec validation checks write ownership: each state table must have exactly one owner writer, and a pipeline where two managers lock the same table for writing is rejected with the error State table <path> is claimed for writing by both .... Read-only consumers (joiners) don't lock the table for writing, so they can share it with the owner writer.
Note
For any state, an empty value corresponds to the absence of a row in the table. If the state is empty after modification, the corresponding row is deleted. By default, emptiness and state clearing are determined automatically (by comparing to the default value); the state type can override this behavior.
State storage
States are stored in YTsaurus dynamic tables. Here's a simple schema example:
|
name |
type |
sort_order |
expression |
|
|
|
|
|
|
|
|
|
|
|
|
|
||
|
|
|
group_by_schema consistency
For correctness and performance, we strongly recommend that you use, as the group_by_schema for a computation, the schema of the first key columns of the dynamic table with states (strictly a prefix of the key columns). This ensures that:
- Only one partition handles events for a single key (correctness).
- One partition handles a limited number of tablets (performance).
Here's an example of a group_by_schema consistent with the state table schema from the example above:
|
name |
type |
expression |
|
|
|
|
|
|
|
StateCache
Flow provides a shared two-level (uncompressed + compressed) LRU cache for states. Configure it at /dynamic_spec/job_tracker/state_cache.
Source: yt/yt/flow/library/cpp/common/spec.h
|
Parameter |
Description |
|
|
Type: NYT::NYTree::TSize |
|
|
Type: NYT::NYTree::TSize |
Implementation in different languages
- C++: the client
TMutableStateKeyClient<TState>(read-write) orTJoinedStateKeyClient<TState>(read-only) returns the accessorTStateAccessor<TState>/TConstStateAccessor<TState>; the same key client works with both internal and external states. Learn more → - Java: YsonStateAccessor, ProtoStateAccessor, ExternalStateAccessor. Learn more →
- Python: ctx.state(), ctx.external_state(), ctx.proto_state(). Learn more →
- Go:
flow.OpenYSONState,flow.OpenProtoState,flow.OpenExternalState, andflow.OpenJoinedExternalStateopen a keyed state accessor. Learn more →