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 (dict or None).
  • 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 (bytes or None).
  • 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 (or None).
  • 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 a PayloadBuilder with the current values.
  • state.set(payload) — save a new value (it accepts a Payload from builder.finish()).
  • state.clear() — delete the state.

Example from Shuffle (EventReducer):

The pattern for working with an external state is:

  1. Get the current state using ctx.external_state(...).
  2. Create a builder with state.to_builder().
  3. Update the required fields with builder.set(...).
  4. 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.

See also