StateAccessor в YTsaurus Flow (Go)

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

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

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

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

Аксессор в Go — это значение, связывающее две вещи: держатель стейта с определённым именем и ключ конкретного входа. Держатели живут в flow.Runtime и доступны напрямую (rt.InternalState(name), rt.ExternalState(name), rt.JoinedExternalState(name)), но обычному компьютейшену они не нужны: аксессор — это и есть удобный вид на стейт одного ключа.

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

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

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

Обратно воркеру уезжает не весь стейт, а дельта: только изменённые записи. Для Raw- и External-стейтов запись выполняется через Set или Clear; YSON- и Proto-стейты записываются изменением значения на месте. Простое чтение ничего не отправляет.

Clear не стирает запись из аксессора, а помечает её удалённой, и удаление доезжает до воркера именно в таком виде. Для компьютейшена разницы нет: стейт, которого запрос не принёс, и стейт, очищенный в этом запросе, одинаково читаются как отсутствующий — компьютейшен видит стейт таким, каким он станет после ответа.

Важно

flow.Runtime и все открытые из него аксессоры принадлежат горутине, обслуживающей запрос, и не рассчитаны на конкурентное использование. Если обработчик распараллеливает работу, стейт следует читать и писать в той же горутине, в которой обработчик был вызван. Правила запуска дочерних горутин описаны в разделе «Горутины в обработчике».

Виды аксессоров

Go SDK предоставляет пять видов аксессоров:

Аксессор Формат Открытие Описание
RawStateAccessor []byte flow.OpenRawState(rt, name, input) Сырые байты без сериализации
YSONState YSON flow.OpenYSONState[T](rt, name, input) Сериализация Go-значения в YSON
ProtoState Protobuf flow.OpenProtoState[T](rt, name, input) Сериализация через Protobuf
ExternalStateAccessor Go-структура (строка таблицы) flow.OpenExternalState(rt, "/name", input) Чтение и запись строки внешней таблицы
JoinedExternalStateAccessor Go-структура flow.OpenJoinedExternalState(rt, "/name", input) Read-only доступ к таблице чужого стейта

Первые три работают с внутренним стейтом, таблицами которого управляет Flow. RawStateAccessor сериализует явные записи; YSONState и ProtoState держат изменяемое значение и сохраняют его автоматически после успешного батча.

ExternalStateAccessor и JoinedExternalStateAccessor работают с внешним стейтом — динамической таблицей, которую пользователь создаёт сам. Различаются они правами: первый доступен компьютейшену, который стейтом владеет, второй — компьютейшену, который его только читает.

API внутренних стейтов

Raw-аксессор использует явные операции чтения и записи:

Метод Тип результата Описание
Get() ([]byte, bool) Копия сохранённых байтов
Or(fallback []byte) []byte Текущее значение или fallback
Set(data []byte) error Сохранить байты; пустые байты удаляют стейт
Clear() error Удалить стейт

YSON- и Proto-стейты изменяются на месте:

Метод YSONState[T] ProtoState[T, PT]
Empty() bool bool
Get() (*T, bool) (PT, bool)
Value() / Or(fallback) *T PT
Set(value) нет метода error
Clear() без результата без результата
ReadOnly() ReadOnlyYSONState[T] ReadOnlyProtoState[T, PT]

Десериализация выполняется при открытии. Изменения значения сохраняются автоматически только после успешного завершения обработчиков батча — см. Изменение значения на месте.

ExternalStateAccessor преобразует строку таблицы в Go-структуру:

Метод Тип результата Описание
ConvertTo(&value) (bool, error) Прочитать строку в структуру; bool отличает отсутствующую строку
ConvertFrom(&value) error Сохранить поля структуры в строку
Clear() error Удалить строку

JoinedExternalStateAccessor предоставляет ConvertTo(&value), но не запись и не очистку. Для динамических схем у обоих аксессоров остаются низкоуровневые Get, Or и Schema, а у владельца — Builder и Set.

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

Аксессор открывается свободной функцией внутри OnMessage, OnTimer или OnVisit:

func (*myFunction) OnMessage(
    ctx context.Context,
    rt flow.Runtime,
    msg flow.ExtendedMessage,
    out flow.OutputCollector,
) error {
    // YSON
    ysonState, err := flow.OpenYSONState[myState](rt, "state-name", msg)

    // Сырые байты
    rawState, err := flow.OpenRawState(rt, "state-name", msg)

    // Protobuf
    protoState, err := flow.OpenProtoState[TMyState](rt, "state-name", msg)

    // Внешний стейт: имя обязательно начинается с "/"
    extState, err := flow.OpenExternalState(rt, "/state-name", msg)

    // Внешний стейт только на чтение
    joined, err := flow.OpenJoinedExternalState(rt, "/reference", msg)
}

func (*myFunction) OnTimer(
    ctx context.Context,
    rt flow.Runtime,
    timer flow.Timer,
    out flow.OutputCollector,
) error {
    // То же самое, но вместо сообщения передаётся таймер
    state, err := flow.OpenYSONState[myState](rt, "state-name", timer)
}

Параметры:

  • rt — flow.Runtime, второй аргумент обработчика.
  • name — имя стейта, объявленное в статической спеке. Для внутренних стейтов — произвольная строка из internal_states; для внешних — ключ из external_state_managers или external_state_joiners, обязательно начинающийся с /.
  • input — вход, к ключу которого привязывается аксессор. Подходит любое значение, реализующее flow.Input: flow.ExtendedMessage, flow.Timer и flow.Visit.

Типовой параметр указывается только у YSON- и Proto-аксессоров и задаёт тип стейта: flow.OpenProtoState[TMyState] выдаёт аксессор над *TMyState. Именно из-за этого функции открытия свободные, а не методы Runtime: собственных типовых параметров у методов в Go нет.

Открытие возвращает ошибку, если имя стейта не подходит:

Ошибка Причина
flow.ErrUnknownState Имя не объявлено в спеке компьютейшена
flow.ErrInvalidStateName Имя внешнего стейта не является абсолютным путём
flow.ErrStateNotRead Запрос не принёс внешнего стейта с таким именем
flow.ErrNoStateSchema Внешний стейт пришёл без схемы своих строк

Примечание

flow.ErrStateNotRead при открытии заджойненного стейта — штатная ситуация, а не сбой: воркер джойнит только те ключи, для которых нашёл строки, поэтому батч, не совпавший ни с одной строкой справочника, приходит без такого стейта. Отличить эту ошибку следует через errors.Is и трактовать как отсутствие данных.

Когда использовать какой тип

Ситуация Рекомендуемый аксессор
Структура из нескольких полей, схема которой меняется вместе с кодом flow.OpenYSONState[T] (YSON)
Произвольные двоичные данные или собственная сериализация flow.OpenRawState (Raw)
Структурированные данные с фиксированной схемой, общей с другими языками flow.OpenProtoState[T] (Protobuf)
Данные, к которым нужен доступ из других систем flow.OpenExternalState (External)
Данные, требующие пользовательской таблицы flow.OpenExternalState (External)
Справочник, который компьютейшен только читает flow.OpenJoinedExternalState (Joined)

См. также

Следующая