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) |
См. также
- Internal State (Go)
- External State (Go)
- Работа со стейтами (Go) — краткий обзор всех типов
- Computation (Go)
- Stateful-вычисления