Working with states in YTsaurus Flow (Python)
Note
This page describes the Python API for working with states. For general state concepts, see the Stateful processing section.
YSON State
The simplest way to work with a state is using the YSON format. You store the state as a Python Map, which is automatically serialized to YSON:
state = ctx.state("state-name", message)
This returns a YsonStateAccessor with the following methods:
get()— get the current value (dictorNone).set(dict)— save the value.clear()— delete the state.get_or_default(dict)— get the current value or return the default value.
Example from WordCount:
Here, the state is tied to the message key, which is defined via group_by_schema in the spec. Each unique key has its own independent state.
Raw State
To store a state as raw bytes:
state = ctx.raw_state("state-name", message)
This returns a RawStateAccessor with the following methods:
get()— get the value (bytesorNone).set(bytes)— save the value.clear()— delete the state.get_or_default(bytes)— get the value or return the default value.
Proto State
To store a state as a Protobuf message:
state_accessor = ctx.proto_state("state-name", message, TJoinState)
This returns a ProtoStateAccessor with the following methods:
get()— deserialize and return the Protobuf object (orNone).set(proto_message)— serialize and save the value.clear()— delete the state.get_or_default(default=None)— get the value or return the default. If you don’t specify a default, it returns an empty instance of the Proto class.
External State
An external state works like a Payload — it gives you dict-like access to fields:
state = ctx.external_state("/state-name", message)
The state name must start with / and match the key in external_state_managers in the static spec. If you call ctx.external_state("state-name", message) without the leading /, it raises a ValueError.
This returns an ExternalStateAccessor, which is also a Payload:
state.get("field")— read a field.state["field"]— dict-like read of a field.state.to_builder()— get aPayloadBuilderwith the current values.state.set(payload)— save a new value (it accepts a Payload frombuilder.finish()).state.clear()— delete the state.
Example from Shuffle (EventReducer):
The pattern for working with an external state is:
- Get the current state using
ctx.external_state(...). - Create a builder with
state.to_builder(). - Update the required fields with
builder.set(...). - Save the changes with
state.set(builder.finish()).
State in timers
The API for working with a state in a timer handler is the same — you pass the timer object instead of message:
def on_timer(self, timer, output, ctx):
state = ctx.external_state("/join-state", timer)
# Read the state
show_time = state.get("show_time")
hit_payload = state.get("hit_payload")
# Clear the state after processing
state.clear()
Example from WaitClickJoin (JoinProcessFunction):
Binding a state to a key
In TTransformCompanionComputation, state binds to the message key defined by group_by_schema in the computation spec. Messages with the same key share one state.
TTransformOrderedSourceCompanionComputation does not support group_by_schema. Its internal state uses the source partition key, so messages from one partition share that state. See Computation (Python) for choosing a SourceComputation class, and Stateful processing for key configuration.
Configuring states in the spec
Internal states (YSON, Raw, Proto) must be declared in the computation parameters in the internal_states section. The state name in your code (the first argument of ctx.state(...)) must match the name declared in the spec.
External states (External) are configured via the external_state_managers section in the computation spec and have their own schema that describes the available fields. The key inside external_state_managers (for example, "/shuffle-state") sets the state name, which must start with /. In the external_state_manager_class_name field, specify the registered manager class (for a typical scenario, use "NYT::NFlow::TSimpleExternalStateManager"). For more details on the spec and available managers, see External State and the C++ documentation.