Internal State в YTsaurus Flow (Go)
Internal State — механизм работы с внутренним состоянием (стейтом), хранящимся во внутренних таблицах Flow. В отличие от External State, пользователю не нужно самостоятельно создавать таблицы — Flow управляет ими автоматически.
Подробнее про аксессоры и общие принципы работы со стейтом: State Accessor (Go).
Общие сведения о stateful-обработке описаны в разделе Stateful processing.
Обзор
Go SDK предоставляет три вида аксессоров для работы с Internal State, различающихся форматом сериализации. Для работы с внешним стейтом используются отдельные аксессоры — External State (Go).
| Аксессор | Формат | Открывается функцией |
|---|---|---|
| YSONState | YSON | flow.OpenYSONState[T] |
| RawStateAccessor | []byte |
flow.OpenRawState |
| ProtoState | Protobuf | flow.OpenProtoState[T] |
RawStateAccessor читает и записывает сырые байты явно через Get, Set и Clear. YSONState и ProtoState предоставляют изменяемое значение: изменения, сделанные в нём, сериализуются автоматически после успешного завершения обработчиков батча — см. Изменение значения на месте.
Каждая из функций flow.OpenXxxState принимает три аргумента:
rt—flow.Runtimeобработчика.name— имя стейта. Имена внутренних стейтов не начинаются с/и должны быть объявлены вparameters.internal_statesспеки компьютейшена, см. Конфигурация в статической спеке.input— вход, к ключу которого привязывается стейт. Подходит любое значение, реализующееflow.Input:flow.ExtendedMessage,flow.Timer,flow.Visit.
Примечание
Аксессор показывает стейт таким, каким он станет после ответа воркеру: стейт, которого не было во входном запросе, и стейт, очищенный в этом же вызове через Clear, читаются одинаково — как отсутствующий.
Изменение значения на месте
Значение, которое возвращают Value(), Get() и Or(), живое: для каждого ключа оно декодируется один раз за батч, все аксессоры этого ключа возвращают одно и то же значение, а изменения, сделанные в нём, записываются в стейт по окончании батча без вызова Set(). Если значение не изменилось, запись не выполняется: стейт, который сериализуется в те же байты, с которыми пришёл, обратно не уезжает. Дефолт из Value() и Or() становится значением стейта и записывается, как после Set(), поэтому его можно сразу изменять:
state, err := flow.OpenYSONState[wordCountState](rt, "word-state", msg)
if err != nil {
return err
}
state.Value().Count++
ProtoState.Set() по-прежнему заменяет сообщение целиком, Clear() удаляет стейт. Пустое Protobuf-сообщение сериализуется в ноль байт, а ноль байт — это отсутствие стейта: дефолт Or(&T{}) не записывается, а сообщение, у которого сняли все поля, и Set(&T{}) стейт удаляют. RawStateAccessor изменяемого значения не отдаёт: Get() возвращает копию байтов, и записываются они только через Set(). Стейты External State на месте тоже не отслеживаются.
У одного ключа за батч только одно декодированное значение, поэтому открытие того же стейта и ключа с другим Go-типом возвращает ошибку. Если обработчик вернул ошибку, изменения на месте не попадают в ответ воркеру. Стейт, который Go сериализует не в те байты, с которыми он пришёл, при первом же чтении перезаписывается канонической записью — один раз на ключ.
Чтобы отследить изменения, прочитанное значение перекодируется по окончании батча. Если компьютейшен только читает стейт, используйте ReadOnly(): такое представление возвращает то же значение, но не отслеживает его, а ReadOnlyYSONState.Value() и ReadOnlyProtoState.Or() не создают стейт. Методов записи у него нет, поэтому запись через него не компилируется. Значение общее с изменяемым стейтом того же ключа: если тот же стейт где-то открыт и на запись, изменения, сделанные через представление, всё же уедут воркеру.
count := state.ReadOnly().Value().Count
YSONState
YSONState[T] хранит стейт как YSON-сериализованное значение типа T. Типом может быть любая структура с тегами yson, а также map, slice или скаляр — всё, что понимает yson.Marshal.
Получение стейта
// Для сообщения
state, err := flow.OpenYSONState[wordCountState](rt, "word-state", msg)
// Для таймера
state, err := flow.OpenYSONState[wordCountState](rt, "word-state", timer)
Десериализация выполняется при открытии. Повторное открытие того же стейта и ключа в пределах запроса возвращает то же изменяемое значение.
Методы
| Метод | Возвращаемый тип | Описание |
|---|---|---|
Empty() |
bool |
Проверить, отсутствует ли значение |
Get() |
(*T, bool) |
Получить изменяемое значение. Второй результат отличает сохранённый стейт от отсутствующего |
Value() |
*T |
Получить изменяемое значение; для отсутствующего стейта создаётся zero value |
Clear() |
— | Удалить значение |
ReadOnly() |
ReadOnlyYSONState[T] |
Представление только на чтение |
Изменения значения сериализуются автоматически после успешного завершения всех обработчиков батча. Если обработчик вернул ошибку, изменения YSON-стейта не попадают в ответ воркеру.
Пример из WordCount
Тип, который пайплайн хранит для одного слова:
type wordCountState struct {
Word string `yson:"word"`
Count int64 `yson:"count"`
}
Обработчик сообщения:
type wordCountMapper struct{}
var _ flow.RowFunction = (*wordCountMapper)(nil)
func (*wordCountMapper) OnMessage(
ctx context.Context,
rt flow.Runtime,
msg flow.ExtendedMessage,
out flow.OutputCollector,
) error {
var input wordMessage
if err := msg.ConvertTo(&input); err != nil {
return err
}
state, err := flow.OpenYSONState[wordCountState](rt, wordStateName, msg)
if err != nil {
return err
}
fresh := state.Empty()
counter := state.Value()
if fresh {
counter.Word = input.Word
}
counter.Count++
return nil
}
Здесь стейт привязан к ключу сообщения. Для нового ключа Empty() возвращает true, а Value() создаёт пустой wordCountState. Присваивания полям сохраняются без отдельного Set.
Тот же стейт открывается и в обработчике таймера. Так устроен URL Downloader: OnMessage накапливает батч, а OnTimer читает его и очищает через Clear.
RawStateAccessor
RawStateAccessor работает с сырыми байтами без сериализации и десериализации. Это аксессор, поверх которого построены остальные два, — берите его, когда формат стейта определяете вы сами.
Получение аксессора
// Для сообщения
state, err := flow.OpenRawState(rt, "raw-state", msg)
// Для таймера
state, err := flow.OpenRawState(rt, "raw-state", timer)
Методы
| Метод | Возвращаемый тип | Описание |
|---|---|---|
Get() |
([]byte, bool) |
Получить сырые байты. Второй результат отличает сохранённый стейт от отсутствующего |
Or(fallback []byte) |
[]byte |
Вернуть текущее значение или fallback, если стейта нет |
Set(data []byte) |
error |
Сохранить сырые байты; пустые байты удаляют стейт |
Clear() |
error |
Удалить стейт для текущего ключа |
Методы Get и Or не возвращают ошибку: десериализовать здесь нечего.
Пример использования
state, err := flow.OpenRawState(rt, "raw-state", msg)
if err != nil {
return err
}
// Чтение сырых данных
if data, ok := state.Get(); ok {
// Обработка сырых данных...
_ = data
}
// Запись сырых данных
if err := state.Set([]byte{0x01, 0x02, 0x03}); err != nil {
return err
}
// Очистка
return state.Clear()
ProtoState
ProtoState сериализует стейт через Protobuf. Тип Protobuf-сообщения указывается в значимой форме, а стейт отдаёт указатель на него: flow.OpenProtoState[TJoinState] возвращает стейт над *TJoinState.
Получение стейта
// Для сообщения
state, err := flow.OpenProtoState[TJoinState](rt, "join-state", msg)
// Для таймера
state, err := flow.OpenProtoState[TJoinState](rt, "join-state", timer)
Десериализация выполняется при открытии. Повторное открытие того же стейта и ключа в пределах запроса возвращает то же изменяемое сообщение.
Методы
| Метод | Возвращаемый тип | Описание |
|---|---|---|
Empty() |
bool |
Проверить, отсутствует ли значение |
Get() |
(*T, bool) |
Получить изменяемое сообщение. Второй результат отличает сохранённый стейт от отсутствующего |
Or(fallback *T) |
*T |
Вернуть текущее сообщение или записать и вернуть fallback, если стейта нет |
Set(value *T) |
error |
Заменить значение стейта целиком; пустое сообщение удаляет стейт |
Clear() |
— | Удалить стейт для текущего ключа |
ReadOnly() |
ReadOnlyProtoState[T, PT] |
Представление только на чтение |
Примечание
В отличие от Python, где get_or_default() без аргументов отдаёт пустой экземпляр Proto-класса, в Go значение по умолчанию задаётся явно — передайте &T{}, если хотите начать с пустого сообщения. Пустое сообщение сериализуется в ноль байт, а ноль байт означают отсутствие стейта: такое сообщение не записывается для ключа, у которого стейта нет, и удаляет стейт у ключа, у которого он есть. Учтите, что сообщение, все поля которого равны нулевым значениям, пустое именно в этом смысле.
Пример использования
state, err := flow.OpenProtoState[TJoinState](rt, "join-state", msg)
if err != nil {
return err
}
window := state.Or(&TJoinState{})
window.ShowTime = showTime
return nil
Конфигурация в статической спеке
Internal State не требует создания внешних таблиц. Стейты автоматически хранятся во внутренних таблицах Flow.
Имена внутренних стейтов должны быть объявлены в секции internal_states параметров компьютейшена в статической спеке:
{
"spec" = {
"computations" = {
"reader" = {
"computation_class_name" = "NYT::NFlow::TSwiftPassthroughOrderedSourceComputation";
"output_stream_ids" = ["words"];
"source_streams" = {
"queue" = {
"source_class_name" = "NYT::NFlow::TQueueSource";
"parameters" = {
"queue_path" = "<cluster=cluster_name>//path/to/queue";
"consumer_path" = "<cluster=cluster_name>//path/to/consumer";
"finite" = false;
};
};
};
"parameters" = {};
};
"mapper" = {
"computation_class_name" = "NYT::NFlow::NCompanion::TTransformCompanionComputation";
"group_by_schema" = [
{"name" = "hash"; "expression" = "farm_hash(word)"; "type" = "uint64"; required = %true;};
{"name" = "word"; "type" = "string";};
];
"input_stream_ids" = ["words"];
"output_stream_ids" = [];
"required_resource_ids" = {
"CompanionManager" = {
"worker" = true;
"controller" = false;
};
};
"parameters" = {
"internal_states" = ["word-state"];
};
};
};
"resources" = {
"CompanionManager" = {
"resource_class_name" = "NYT::NFlow::NCompanion::TCompanionManager";
"parameters" = {
};
"dependencies" = {};
};
};
};
"dynamic_spec" = {
"computations" = {
"reader" = {
"parameters" = {
};
};
"mapper" = {
"parameters" = {
};
};
};
};
}
Имя стейта в коде (второй аргумент flow.OpenYSONState, flow.OpenRawState или flow.OpenProtoState) должно совпадать с именем, объявленным в internal_states.
Важно
Если имя стейта не объявлено в internal_states, функция открытия вернёт ошибку, обёртывающую flow.ErrUnknownState; текст ошибки перечисляет объявленные имена. Возвращённая из обработчика ошибка прекращает обработку всего батча — воркер повторит запрос целиком.