Работа со стейтами в YTsaurus Flow (Java)
Примечание
На этой странице описаны детали для Java и Kotlin при работе со стейтами. Общие концепции описаны в разделе Stateful-вычисления.
Java SDK Flow (Java и Kotlin) предоставляет несколько типов стейт-аксессоров для работы с состоянием. Наиболее часто используются:
- YsonStateAccessor — для YSON-стейтов, хранящихся во внутренних таблицах Flow.
- ExternalStateAccessor — для внешних стейтов, хранящихся в отдельных динамических таблицах YTsaurus.
Также доступны ProtoStateAccessor, DefaultStateAccessor, RawStateAccessor и NoOpStateAccessor. Подробнее о всех типах — в разделе Internal State.
YsonStateAccessor
YsonStateAccessor<T> предоставляет типизированный доступ к YSON-стейту, привязанному к ключу сообщения. Получить аксессор можно через RuntimeContext:
StateAccessor<T> stateAccessor = ctx.getYsonStateAccessor("state-name", message, StateClass.class);
val stateAccessor: StateAccessor<T> = ctx.getYsonStateAccessor("state-name", message, StateClass::class.java)
Параметры:
"state-name"— имя стейта, должно совпадать с именем, зарегистрированным в спеке компьютейшена (parameters/internal_states).message— текущее сообщение, из которого извлекается ключ группировки.StateClass.class— Java-класс для сериализации/десериализации стейта.
Методы StateAccessor
| Метод | Описание |
|---|---|
get() |
Получить текущее значение стейта (T, null если значения нет) |
set(T value) |
Установить новое значение стейта |
getOrDefault(T defaultValue) |
Получить значение или вернуть значение по умолчанию |
clear() |
Удалить стейт для текущего ключа |
Пример: WordCountMapper
@Component
public class WordCountMapper implements RowFunction {
@Override
public void onMessage(ExtendedMessage message, OutputCollector output, RuntimeContext ctx) {
Word input = message.getPayload();
StateAccessor<WordCountState> stateAccessor =
ctx.getYsonStateAccessor("word-state", message, WordCountState.class);
var state = stateAccessor.getOrDefault(new WordCountState(input.getWord(), 0));
state.setCount(state.getCount() + 1);
stateAccessor.set(state);
}
}
@Component
class WordCountMapper : RowFunction {
override fun onMessage(message: ExtendedMessage, output: OutputCollector, ctx: RuntimeContext) {
val input: Word = message.getPayload()
val stateAccessor = ctx.getYsonStateAccessor("word-state", message, WordCountState::class.java)
val state = stateAccessor.getOrDefault(WordCountState(input.word, 0))
state.count = state.count + 1
stateAccessor.set(state)
}
}
В этом примере:
- Из сообщения извлекается объект
Wordс полемword. - Получается аксессор для стейта
"word-state", привязанного к ключу текущего сообщения. - Если стейт для данного ключа отсутствует, создается новый объект
WordCountStateс начальным значением счетчика 0. - Значение счетчика увеличивается и стейт обновляется.
Класс стейта должен быть аннотирован @YTreeObject для сериализации в YSON:
@YTreeObject
public class WordCountState {
private String word;
private long count;
public WordCountState() {}
public WordCountState(String word, long count) {
this.word = word;
this.count = count;
}
// getters и setters
public String getWord() { return word; }
public void setWord(String word) { this.word = word; }
public long getCount() { return count; }
public void setCount(long count) { this.count = count; }
}
@YTreeObject
class WordCountState {
var word: String = ""
var count: Long = 0
constructor()
constructor(word: String, count: Long) {
this.word = word
this.count = count
}
}
ExternalStateAccessor
ExternalStateAccessor предоставляет доступ к внешнему стейту, хранящемуся в отдельной динамической таблице YTsaurus. Внешний стейт описывается константой ExternalStateDescriptor, создаваемой через StateDescriptors.external(...):
ExternalStateAccessor externalStateAccessor = ctx.getExternalStateAccessor("state-name", message);
val externalStateAccessor = ctx.getExternalStateAccessor("state-name", message)
Методы ExternalStateAccessor
| Метод | Описание |
|---|---|
get() |
Получить текущее значение стейта (Payload, null если значения нет) |
getOrDefault() |
Получить значение или пустой Payload |
set(Payload value) |
Установить новое значение стейта |
clear() |
Удалить стейт для текущего ключа |
Payload — нетипизированный контейнер с доступом к полям по имени. Для модификации используется PayloadBuilder.
Пример: EventReducer
public class EventReducer implements RowFunction {
@Override
public void onMessage(ExtendedMessage message, OutputCollector output, RuntimeContext ctx) {
ExternalStateAccessor externalStateAccessor =
ctx.getExternalStateAccessor("shuffle-state", message);
Payload state = externalStateAccessor.getOrDefault();
PayloadBuilder stateBuilder = state.toBuilder();
if (state.get("count", Long.class) == null) {
stateBuilder.set("count", 1L);
} else {
stateBuilder.set("count", state.get("count", Long.class) + 1);
}
externalStateAccessor.set(stateBuilder.finish());
}
}
class EventReducer : RowFunction {
override fun onMessage(message: ExtendedMessage, output: OutputCollector, ctx: RuntimeContext) {
val externalStateAccessor = ctx.getExternalStateAccessor("shuffle-state", message)
val state = externalStateAccessor.getOrDefault()
val stateBuilder = state.toBuilder()
if (state.get("count", Long::class.java) == null) {
stateBuilder.set("count", 1L)
} else {
stateBuilder.set("count", state.get("count", Long::class.java) + 1)
}
externalStateAccessor.set(stateBuilder.finish())
}
}
В этом примере:
- Объявляется дескриптор
SHUFFLE_STATEдля внешнего стейта"/shuffle-state". - Через
ctx.getState(SHUFFLE_STATE, message)получается аксессор для ключа сообщения. - Текущее значение стейта извлекается как
Payload. - С помощью
PayloadBuilderсоздается обновленная версия стейта с увеличенным счетчиком. - Обновленный стейт сохраняется обратно.
Стейт в таймерах
При обработке таймеров стейт доступен через объект timer, который содержит ключ группировки:
@Override
public void onTimer(Timer timer, OutputCollector output, RuntimeContext ctx) {
ExternalStateAccessor stateAccessor =
ctx.getExternalStateAccessor("join-state", timer);
Payload joinState = Objects.requireNonNull(stateAccessor.get());
// обработка стейта и генерация выходных сообщений
var messageBuilder = ctx.createMessageBuilder("output_stream");
messageBuilder.set("hit_id", joinState.get("hit_id", String.class));
// ... заполнение остальных полей ...
output.addMessage(messageBuilder.finish());
// очистка стейта после обработки
stateAccessor.clear();
}
override fun onTimer(timer: Timer, output: OutputCollector, ctx: RuntimeContext) {
val stateAccessor = ctx.getExternalStateAccessor("join-state", timer)
val joinState = stateAccessor.get()!!
// обработка стейта и генерация выходных сообщений
val messageBuilder = ctx.createMessageBuilder("output_stream")
messageBuilder.set("hit_id", joinState.get("hit_id", String::class.java))
// ... заполнение остальных полей ...
output.addMessage(messageBuilder.finish())
// очистка стейта после обработки
stateAccessor.clear()
}
Метод clear() удаляет стейт для данного ключа. Это важно делать после закрытия окна или финализации обработки, чтобы не накапливать устаревшие данные.
Привязка стейта к ключу
В TTransformCompanionComputation ключ, по которому осуществляется доступ к стейту, определяется полем group_by_schema в спеке компьютейшена. Стейт-аксессор автоматически извлекает ключ из переданного сообщения или таймера.
В TTransformOrderedSourceCompanionComputation поле group_by_schema не поддерживается. Ключом внутреннего стейта служит ключ партиции источника, поэтому все сообщения одной партиции разделяют стейт. Подробнее о выборе класса для SourceComputation см. в разделе Computation (Java).
Подробнее про group_by_schema и его влияние на работу стейтов — в разделе Stateful-вычисления.
Конфигурация стейта в спеке
YSON-стейт (internal_states)
YSON-стейты регистрируются в спеке компьютейшена через parameters/internal_states:
"computations" = {
"mapper" = {
"computation_class_name" = "NYT::NFlow::NCompanion::TTransformCompanionComputation";
"parameters" = {
"internal_states" = ["word-state"];
};
};
};
Внешний стейт (external state)
Внешний стейт регистрируется в секции external_state_managers Computation. Имя должно начинаться с / и совпадать со значением, переданным в StateDescriptors.external(...):
"computations" = {
"reducer" = {
"computation_class_name" = "NYT::NFlow::NCompanion::TTransformCompanionComputation";
"external_state_managers" = {
"/shuffle-state" = {
"external_state_manager_class_name" = "NYT::NFlow::TSimpleExternalStateManager";
"parameters" = {
"path" = "//path/to/state/table";
};
};
};
"parameters" = {};
};
};
Поле external_state_manager_class_name задаёт зарегистрированный класс external state manager — для типового сценария это "NYT::NFlow::TSimpleExternalStateManager". Подробнее про доступные менеджеры см. в C++ документации.
Таблицу стейта необходимо создать заблаговременно.