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; текст ошибки перечисляет объявленные имена. Возвращённая из обработчика ошибка прекращает обработку всего батча — воркер повторит запрос целиком.

См. также

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