StateAccessor in YTsaurus Flow (Java)

StateAccessor is an interface for reading, modifying, and deleting state values. For general information about stateful processing, see the Stateful processing section.

How it works

The state in Flow is stored in sorted dynamic tables.
If you use external state, you create this table. If you use internal state, Flow creates and manages these tables automatically.

For simplicity, the following description focuses on an example with external state.

You can think of each row in the state table as having two parts: key columns and value columns:

hash word count system attributesKey columnslengthValue columns

For TTransformCompanionComputation, the key columns in the state table match the group_by_schema of the computation. For the internal state of TTransformOrderedSourceCompanionComputation, the source partition key is used instead; this SourceComputation does not support group_by_schema.

The value columns are available for reading and modifying through StateAccessor. The format in which you can read and modify these values in Java code depends on the StateAccessor implementation.

Reading and writing data

The worker handles direct operations on the table, including reading, writing, and deleting data. When the worker receives the next batch of messages, it loads the state values for all keys in the batch and sends them to the companion along with the messages and timers. For more details, see the interaction schema.

You write new values to the state table as a transaction within an epoch.

StateAccessor interface

public interface StateAccessor<T> {
    /** Get the state value. */
    @Nullable
    T get();

    /** Get the state value or a default value. */
    default T getOrDefault(T defaultValue);

    /** Set the state value. */
    void set(T value);

    /** Clear or delete the state for the key. */
    void clear();

    /** Get the state class. */
    Class<T> getStateClass();

    /** Get a read-only view of the accessor. */
    default StateAccessor<T> readOnly();
}
interface StateAccessor<T> {
    /** Get the state value. */
    fun get(): T?

    /** Get the state value or a default value. */
    fun getOrDefault(defaultValue: T): T

    /** Set the state value. */
    fun set(value: T)

    /** Clear or delete the state for the key. */
    fun clear()

    /** Get the state class. */
    fun getStateClass(): Class<T>

    /** Get a read-only view of the accessor. */
    fun readOnly(): StateAccessor<T>
}

The value of an internal state returned by get() and getOrDefault() is live: the changes made to it are written without a set() call, and readOnly() returns an untracked view — see Changing the value in place. A default returned by getOrDefault() is not written to YTsaurus: it becomes the state value only if the computation changes it.