StateAccessor в YTsaurus Flow (Java)

StateAccessor — интерфейс для чтения, модификации и удаления значений стейта.
Общие сведения о stateful-обработке описаны в разделе Stateful processing.

Принцип работы

Стейт во Flow хранится в сортированных динамических таблицах.
В случае внешнего стейта эта таблица создаётся пользователем, в случае внутреннего стейта эти таблицы создаются и управляются Flow автоматически.

Далее для простоты описания будем рассматривать пример внешнего стейта.

Каждую строку в таблице стейта можно условно разделить на ключевые колонки и колонки значений:

hash word count system attributesKey columnslengthValue columns

Для TTransformCompanionComputation ключевые колонки в таблице стейта совпадают с group_by_schema компьютейшена. Для внутреннего стейта TTransformOrderedSourceCompanionComputation ключом служит ключ партиции источника: group_by_schema в таком SourceComputation не поддерживается.

Колонки значений будут доступны для чтения и модификации через StateAccessor. Формат, в котором эти значения будут доступны для чтения и модификации в Java-коде, зависит от реализации StateAccessor.

Чтение и запись данных

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

Запись новых значений в таблицу стейта осуществляется транзакционно в рамках эпохи.

Интерфейс StateAccessor

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

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

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

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

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

    /** Получить read-only представление аксессора. */
    default StateAccessor<T> readOnly();
}
interface StateAccessor<T> {
    /** Получить значение стейта. */
    fun get(): T?

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

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

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

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

    /** Получить read-only представление аксессора. */
    fun readOnly(): StateAccessor<T>
}

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

Следующая