Internal State в YTsaurus Flow (Java)

Internal State — механизм работы с внутренним состоянием (стейтом), хранящимся во внутренних таблицах Flow. В отличие от External State, пользователю не нужно самостоятельно создавать таблицы — Flow управляет ими автоматически.

Подробнее про StateAccessor и работу со стейтом: State Accessor.

Общие сведения о stateful-обработке описаны в разделе Stateful processing.

Обзор

Java SDK Flow (Java и Kotlin) предоставляет несколько видов аксессоров состояния для работы с Internal State, различающихся форматом сериализации:

Accessor Формат Описание
YsonStateAccessor YSON Сериализация через @YTreeObject аннотации
ProtoStateAccessor Protobuf Сериализация через Protobuf
DefaultStateAccessor Произвольный Пользовательские serializer/deserializer
RawStateAccessor byte[] Без сериализации (сырые байты)
NoOpStateAccessor - Хранит только факт наличия стейта

Все аксессоры реализуют общий интерфейс StateAccessor<T>.

Интерфейс StateAccessor

public interface StateAccessor<T> {
    /** Получить значение стейта. */
    @Nullable
    T get();

    /** Получить значение стейта или дефолтное значение. */
    default T getOrDefault(T defaultValue);

    /** Установить значение стейта. */
    void set(T value);

    /** Очистить/удалить стейт для ключа. */
    void clear();

    /** Получить класс стейта. */
    Class<T> getStateClass();
}
interface StateAccessor<T> {
    /** Получить значение стейта. */
    fun get(): T?

    /** Получить значение стейта или дефолтное значение. */
    fun getOrDefault(defaultValue: T): T

    /** Установить значение стейта. */
    fun set(value: T)

    /** Очистить/удалить стейт для ключа. */
    fun clear()

    /** Получить класс стейта. */
    fun getStateClass(): Class<T>
}

Изменение значения на месте

Значение, которое возвращают get() и getOrDefault(), живое: для каждого ключа оно декодируется один раз за батч, все аксессоры этого ключа возвращают один и тот же объект, а изменения, сделанные в нём, записываются в стейт по окончании батча без вызова set(). Если значение не изменилось, запись не выполняется. Дефолт из getOrDefault() привязывается к ключу, но не записывается: его можно сразу изменять, и значением стейта он становится только после изменения — нетронутый дефолт не создаёт строку в таблице стейтов:

ctx.getState(COUNTER, message).getOrDefault(new CounterState()).count += 1;
ctx.getState(COUNTER, message).getOrDefault(CounterState()).count += 1

set() по-прежнему заменяет значение целиком, clear() удаляет стейт. Protobuf-объекты неизменяемы, поэтому для них новое значение задаётся через set().

Чтобы отследить изменения, прочитанное значение один раз перекодируется по окончании батча — даже если вычисление стейт только просматривает. Убрать эту работу позволяет readOnly(): аксессор возвращает то же значение, но не отслеживает его, поэтому по окончании батча стейт не перекодируется и не отправляется воркеру. Используйте его везде, где стейт только читают. getOrDefault() у такого аксессора не создаёт стейт, а set() и clear() бросают исключение:

long count = ctx.getState(COUNTER, message).readOnly().getOrDefault().count;
val count = ctx.getState(COUNTER, message).readOnly().getOrDefault().count

То же самое можно объявить один раз на дескрипторе: InternalStateDescriptor.readOnly() возвращает дескриптор того же стейта, аксессоры которого доступны только на чтение, — тогда вызывать readOnly() на каждом месте обращения не нужно. Такой дескриптор объявляют один раз константой рядом с исходным:

private static final InternalStateDescriptor<CounterState> COUNTER_READ_ONLY = COUNTER.readOnly();

long count = ctx.getState(COUNTER_READ_ONLY, message).getOrDefault().count;
private val COUNTER_READ_ONLY: InternalStateDescriptor<CounterState> = COUNTER.readOnly()

val count = ctx.getState(COUNTER_READ_ONLY, message).getOrDefault().count

Внутренний стейт живёт в пределах одного компьютейшена: воркер отдаёт только те имена, которые перечислены в internal_states его параметров. Поэтому read-only-дескриптор читает стейт того компьютейшена, в котором используется: прочитать им стейт, который пишет другой компьютейшен, нельзя — для этого есть joiner'ы External State.

Read-only — это дисциплина обращения, а не свойство стейта: отслеживание живёт на самом стейте. Если в этом же батче значение для того же ключа уже читали через пишущий аксессор, стейт уже отслеживается, и сделанное после этого изменение на месте дойдёт до воркера, каким бы аксессором значение ни было получено. Дескриптор гарантирует «этот аксессор не пишет», а не «этот стейт не отслеживается».

YsonStateAccessor

YsonStateAccessor использует YSON-сериализацию. Класс стейта должен быть аннотирован @YTreeObject.

Получение аксессора

// Для сообщения
YsonStateAccessor<MyState> stateAccessor =
        ctx.getYsonStateAccessor("state-name", message, MyState.class);

// Для таймера
YsonStateAccessor<MyState> stateAccessor =
        ctx.getYsonStateAccessor("state-name", timer, MyState.class);
// Для сообщения
val stateAccessor: YsonStateAccessor<MyState> =
        ctx.getYsonStateAccessor("state-name", message, MyState::class.java)

// Для таймера
val stateAccessor: YsonStateAccessor<MyState> =
        ctx.getYsonStateAccessor("state-name", timer, MyState::class.java)

Пример класса стейта

import ru.yandex.inside.yt.kosher.impl.ytree.object.annotation.YTreeObject;
import ru.yandex.inside.yt.kosher.impl.ytree.object.annotation.YTreeField;

@YTreeObject
public class CounterState {
    @YTreeField(key = "count")
    private long count;

    @YTreeField(key = "last_update")
    private long lastUpdate;

    // Конструктор по умолчанию обязателен
    public CounterState() {}

    // Геттеры и сеттеры...
    public long getCount() { return count; }
    public void setCount(long count) { this.count = count; }
    public long getLastUpdate() { return lastUpdate; }
    public void setLastUpdate(long lastUpdate) { this.lastUpdate = lastUpdate; }
}
import ru.yandex.inside.yt.kosher.impl.ytree.object.annotation.YTreeObject
import ru.yandex.inside.yt.kosher.impl.ytree.object.annotation.YTreeField

@YTreeObject
class CounterState {
    @YTreeField(key = "count")
    var count: Long = 0

    @YTreeField(key = "last_update")
    var lastUpdate: Long = 0

    // Конструктор по умолчанию обязателен
    constructor()
}

Пример использования

public class CounterFunction implements RowFunction {
    @Override
    public void onMessage(ExtendedMessage message, OutputCollector output, RuntimeContext ctx) {
        YsonStateAccessor<CounterState> stateAccessor =
                ctx.getYsonStateAccessor("counter", message, CounterState.class);

        // Получение текущего стейта или создание нового
        CounterState state = stateAccessor.getOrDefault(new CounterState());

        // Модификация стейта
        state.setCount(state.getCount() + 1);
        state.setLastUpdate(message.getEventTimestamp());
    }
}
class CounterFunction : RowFunction {
    override fun onMessage(message: ExtendedMessage, output: OutputCollector, ctx: RuntimeContext) {
        val stateAccessor: YsonStateAccessor<CounterState> =
                ctx.getYsonStateAccessor("counter", message, CounterState::class.java)

        // Получение текущего стейта или создание нового
        val state: CounterState = stateAccessor.getOrDefault(CounterState())

        // Модификация стейта
        state.count = state.count + 1
        state.lastUpdate = message.getEventTimestamp()
    }
}

ProtoStateAccessor

Исходный код

ProtoStateAccessor использует Protobuf-сериализацию. Класс стейта должен наследовать com.google.protobuf.MessageLite.

Получение аксессора

// Для сообщения
ProtoStateAccessor<MyProtoState> stateAccessor =
        ctx.getProtoStateAccessor("state-name", message, MyProtoState.class);

// Для таймера
ProtoStateAccessor<MyProtoState> stateAccessor =
        ctx.getProtoStateAccessor("state-name", timer, MyProtoState.class);
// Для сообщения
val stateAccessor: ProtoStateAccessor<MyProtoState> =
        ctx.getProtoStateAccessor("state-name", message, MyProtoState::class.java)

// Для таймера
val stateAccessor: ProtoStateAccessor<MyProtoState> =
        ctx.getProtoStateAccessor("state-name", timer, MyProtoState::class.java)

Метод getOrDefault

ProtoStateAccessor предоставляет метод getOrDefault() без параметров, который возвращает дефолтный Protobuf-объект:

ProtoStateAccessor<MyProtoState> stateAccessor =
        ctx.getProtoStateAccessor("state-name", message, MyProtoState.class);

// Получение стейта или дефолтного Protobuf-объекта
MyProtoState state = stateAccessor.getOrDefault();
val stateAccessor: ProtoStateAccessor<MyProtoState> =
        ctx.getProtoStateAccessor("state-name", message, MyProtoState::class.java)

// Получение стейта или дефолтного Protobuf-объекта
val state: MyProtoState = stateAccessor.getOrDefault()

Пример использования

public class ProtoCounterFunction implements RowFunction {
    @Override
    public void onMessage(ExtendedMessage message, OutputCollector output, RuntimeContext ctx) {
        ProtoStateAccessor<CounterProto> stateAccessor =
                ctx.getProtoStateAccessor("counter", message, CounterProto.class);

        CounterProto state = stateAccessor.getOrDefault();

        // Модификация через Protobuf builder
        CounterProto updatedState = state.toBuilder()
                .setCount(state.getCount() + 1)
                .setLastUpdate(message.getEventTimestamp())
                .build();

        stateAccessor.set(updatedState);
    }
}
class ProtoCounterFunction : RowFunction {
    override fun onMessage(message: ExtendedMessage, output: OutputCollector, ctx: RuntimeContext) {
        val stateAccessor: ProtoStateAccessor<CounterProto> =
                ctx.getProtoStateAccessor("counter", message, CounterProto::class.java)

        val state: CounterProto = stateAccessor.getOrDefault()

        // Модификация через Protobuf builder
        val updatedState: CounterProto = state.toBuilder()
                .setCount(state.getCount() + 1)
                .setLastUpdate(message.getEventTimestamp())
                .build()

        stateAccessor.set(updatedState)
    }
}

DefaultStateAccessor

Исходный код

DefaultStateAccessor позволяет использовать произвольные функции сериализации и десериализации.

Получение аксессора

// Для сообщения
DefaultStateAccessor<MyState> stateAccessor = ctx.getStateAccessor(
        "state-name",
        message,
        MyState.class,
        state -> serialize(state),      // Function<MyState, byte[]>
        bytes -> deserialize(bytes)     // Function<byte[], MyState>
);

// Для таймера
DefaultStateAccessor<MyState> stateAccessor = ctx.getStateAccessor(
        "state-name",
        timer,
        MyState.class,
        state -> serialize(state),
        bytes -> deserialize(bytes)
);
// Для сообщения
val stateAccessor: DefaultStateAccessor<MyState> = ctx.getStateAccessor(
        "state-name",
        message,
        MyState::class.java,
        { state -> serialize(state) },      // (MyState) -> ByteArray
        { bytes -> deserialize(bytes) }     // (ByteArray) -> MyState
)

// Для таймера
val stateAccessor: DefaultStateAccessor<MyState> = ctx.getStateAccessor(
        "state-name",
        timer,
        MyState::class.java,
        { state -> serialize(state) },
        { bytes -> deserialize(bytes) }
)

Пример с Jackson

public class JsonCounterFunction implements RowFunction {
    private static final ObjectMapper mapper = new ObjectMapper();

    @Override
    public void onMessage(ExtendedMessage message, OutputCollector output, RuntimeContext ctx) {
        DefaultStateAccessor<CounterState> stateAccessor = ctx.getStateAccessor(
                "counter",
                message,
                CounterState.class,
                state -> {
                    try { return mapper.writeValueAsBytes(state); }
                    catch (Exception e) { throw new RuntimeException(e); }
                },
                bytes -> {
                    try { return mapper.readValue(bytes, CounterState.class); }
                    catch (Exception e) { throw new RuntimeException(e); }
                }
        );

        CounterState state = stateAccessor.getOrDefault(new CounterState());
        state.setCount(state.getCount() + 1);
    }
}
class JsonCounterFunction : RowFunction {
    companion object {
        private val mapper = ObjectMapper()
    }

    override fun onMessage(message: ExtendedMessage, output: OutputCollector, ctx: RuntimeContext) {
        val stateAccessor: DefaultStateAccessor<CounterState> = ctx.getStateAccessor(
                "counter",
                message,
                CounterState::class.java,
                { state ->
                    try { mapper.writeValueAsBytes(state) }
                    catch (e: Exception) { throw RuntimeException(e) }
                },
                { bytes ->
                    try { mapper.readValue(bytes, CounterState::class.java) }
                    catch (e: Exception) { throw RuntimeException(e) }
                }
        )

        val state: CounterState = stateAccessor.getOrDefault(CounterState())
        state.count = state.count + 1
    }
}

RawStateAccessor

Исходный код

RawStateAccessor работает с сырыми байтами без сериализации/десериализации.

Получение аксессора

RawStateAccessor stateAccessor = ctx.getRawStateAccessor("state-name", message);
val stateAccessor: RawStateAccessor = ctx.getRawStateAccessor("state-name", message)

Пример использования

RawStateAccessor stateAccessor = ctx.getRawStateAccessor("raw-state", message);

byte[] data = stateAccessor.get();
if (data != null) {
    // Обработка сырых данных...
}

// Запись сырых данных
stateAccessor.set(new byte[]{0x01, 0x02, 0x03});

// Очистка
stateAccessor.clear();
val stateAccessor: RawStateAccessor = ctx.getRawStateAccessor("raw-state", message)

val data: ByteArray? = stateAccessor.get()
if (data != null) {
    // Обработка сырых данных...
}

// Запись сырых данных
stateAccessor.set(byteArrayOf(0x01, 0x02, 0x03))

// Очистка
stateAccessor.clear()

NoOpStateAccessor

Исходный код

NoOpStateAccessor хранит только факт наличия стейта для ключа, без какого-либо payload. Полезен для отслеживания уже обработанных ключей (дедупликация).

Получение аксессора

NoOpStateAccessor stateAccessor = ctx.getNoOpStateAccessor("seen-keys", message);
val stateAccessor: NoOpStateAccessor = ctx.getNoOpStateAccessor("seen-keys", message)

Пример использования

public class DeduplicationFunction implements RowFunction {
    @Override
    public void onMessage(ExtendedMessage message, OutputCollector output, RuntimeContext ctx) {
        NoOpStateAccessor stateAccessor = ctx.getNoOpStateAccessor("seen-keys", message);

        // Проверяем, был ли ключ уже обработан
        if (stateAccessor.get() != null) {
            // Ключ уже обработан, пропускаем
            return;
        }

        // Отмечаем ключ как обработанный
        stateAccessor.touch();

        // Обработка сообщения...
        output.addMessage(new Message("output", message.getPayload()));
    }
}
class DeduplicationFunction : RowFunction {
    override fun onMessage(message: ExtendedMessage, output: OutputCollector, ctx: RuntimeContext) {
        val stateAccessor: NoOpStateAccessor = ctx.getNoOpStateAccessor("seen-keys", message)

        // Проверяем, был ли ключ уже обработан
        if (stateAccessor.get() != null) {
            // Ключ уже обработан, пропускаем
            return
        }

        // Отмечаем ключ как обработанный
        stateAccessor.touch()

        // Обработка сообщения...
        output.addMessage(Message("output", message.getPayload()))
    }
}

Конфигурация в статической спеке

Internal State не требует создания внешних таблиц. Стейты автоматически хранятся во внутренних таблицах Flow (states и partition_states).

Имена внутренних стейтов должны быть объявлены в секции internal_states параметров компьютейшена в статической спеке:

"computations" = {
    "counter" = {
        "computation_class_name" = "NYT::NFlow::NCompanion::TTransformCompanionComputation";
        "group_by_schema" = [
            {"name" = "hash"; "expression" = "farm_hash(key)"; "type" = "uint64"};
            {"name" = "key"; "type" = "string"};
        ];
        "input_stream_ids" = ["input"];
        "output_stream_ids" = ["output"];
        "parameters" = {
            "internal_states" = ["counter"];
        };
    };
};

Имя стейта в коде (первый аргумент ctx.getYsonStateAccessor(...), ctx.getProtoStateAccessor(...) и т.д.) должно совпадать с именем, объявленным в internal_states.

Предыдущая
Следующая