Computation в YTsaurus Flow (Java)
Примечание
На этой странице описаны детали для Java и Kotlin при работе с компьютейшенами. Общие концепции описаны в разделе Computation.
Типы Computation
Во Flow есть два вида Computation: Swift и Transform. От их выбора зависит способ обеспечения exactly-once гарантий и то, какие преобразования возможно реализовать с их применением.
| Тип | Способ обеспечения гарантий | Применение |
|---|---|---|
Swift |
Код преобразования детерминирован, при необходимости будет вызываться повторно | Stateless преобразования |
Transform |
Результат работы обязательно сохраняется в YTsaurus, поэтому нет требований на какую-либо детерминированность преобразований | Stateful преобразования |
Подробнее про гарантии обработки — в разделе Гарантии обработки.
Для пайплайнов на Java и Kotlin выбор между Swift или Transform осуществляется через указание computation_class_name в статической спеке:
NYT::NFlow::NCompanion::TTransformCompanionComputation— дляTransform.NYT::NFlow::NCompanion::TSwiftMapCompanionComputation— дляSwift.NYT::NFlow::NCompanion::TTransformOrderedSourceCompanionComputation— дляTransform-сорса.NYT::NFlow::NCompanion::TSwiftOrderedSourceCompanionComputation— дляSwift-сорса.
Создание Computation
В коде на Java и Kotlin Computation создаётся через Computation.Builder и регистрируется в PipelineContext.
var join = Computation.builder()
.setComputationId("join")
.setProcessFunction(new JoinProcessFunction())
.build();
val join = Computation.builder()
.setComputationId("join")
.setProcessFunction(JoinProcessFunction())
.build()
В статической спеке создаётся Computation с таким же id (в данном примере join):
"join" = {
"computation_class_name" = "NYT::NFlow::NCompanion::TTransformCompanionComputation";
"group_by_schema" = [
...
];
"input_stream_ids" = [...];
"output_stream_ids" = [...];
"parameters" = {
...
};
"timers" = {};
};
Подробнее про спеку в разделе Spec, DynamicSpec и Config.
Важно
processFunction обязателен (null недопустим): компьютейшены без бизнес-логики в Java не регистрируют. Если нужен passthrough — не регистрируйте компьютейшен в Java вовсе, а в статической спеке укажите C++-класс passthrough в computation_class_name (см. Passthrough Computation).
SourceComputation
SourceComputation — вершина в графе пайплайна, осуществляющая чтение данных из внешних источников. На стороне воркера ей соответствует TSwiftOrderedSourceComputation или TTransformOrderedSourceComputation.
В Java SourceComputation расширяет Computation. Как и у Computation, параметр processFunction обязателен.
В статической спеке для детерминированной обработки без пользовательского стейта указывают TSwiftOrderedSourceCompanionComputation. Если SourceComputation использует внутренний стейт или недетерминированную логику, указывают TTransformOrderedSourceCompanionComputation: воркер материализует выход и фиксирует его вместе со стейтом и смещением источника. Ключ внутреннего стейта в таком компьютейшене — ключ партиции источника.
Параметры
| Параметр | Обязательный | Описание |
|---|---|---|
computationId |
Да | Уникальный идентификатор |
processFunction |
Да | Функция обработки сообщений |
Создание SourceComputation
var reader = SourceComputation.builder()
.setComputationId("hit_reader")
.setProcessFunction(new HitParsingFunction())
.build();
val reader = SourceComputation.builder()
.setComputationId("hit_reader")
.setProcessFunction(HitParsingFunction())
.build()
Для passthrough Source не используйте Java — укажите в спеке NYT::NFlow::TSwiftPassthroughOrderedSourceComputation в computation_class_name и оставьте компьютейшен незарегистрированным в Java-компаньоне. Подробнее — Passthrough Computation.
Взаимодействие с Worker
При инициализации Worker запрашивает у Java-компаньона информацию о зарегистрированных объектах Computation и SourceComputation. Source-компьютейшен на стороне воркера отправляет входные сообщения в Java-компаньон, который применяет к ним ProcessFunction и возвращает результат.
Process Function
Бизнес-логика обработки данных реализуется через Process Function. Для реализации необходимо выбрать один из двух интерфейсов: RowFunction или BatchFunction.
Примечание
Использование RowFunction или BatchFunction — исключительно вопрос бизнес-логики. RowFunction не добавляет накладных расходов на обработку данных относительно использования BatchFunction благодаря тому, что Flow внутри себя осуществляет передачу данных батчами.
RowFunction
RowFunction получает сообщения и таймеры по одному. Интерфейс предоставляет два метода:
onMessage(ExtendedMessage message, OutputCollector output, RuntimeContext ctx)— вызывается для каждого входного сообщения.onTimer(Timer timer, OutputCollector output, RuntimeContext ctx)— вызывается при срабатывании таймера.
Пример stateless-функции
public class X2Mapper implements RowFunction {
@Override
public void onMessage(ExtendedMessage message, OutputCollector output, RuntimeContext ctx) {
var messageBuilder = ctx.createMessageBuilder("x2_numbers"); //1
Long number = message.get("number", Long.class); //2
messageBuilder.set("number_x2", number * 2); //3
output.addMessage(messageBuilder.finish()); //4
}
}
class X2Mapper : RowFunction {
override fun onMessage(message: ExtendedMessage, output: OutputCollector, ctx: RuntimeContext) {
val messageBuilder = ctx.createMessageBuilder("x2_numbers") //1
val number: Long? = message.get("number", Long::class.java) //2
messageBuilder.set("number_x2", number!! * 2) //3
output.addMessage(messageBuilder.finish()) //4
}
}
Разберем построчно:
-
ctx.createMessageBuilder("x2_numbers")— создаетсяMessageBuilderдля output-стрима с id =x2_numbers. Стрим с таким идентификатором должен присутствовать в спискеoutput_stream_idsв статической спеке компьютейшена. -
message.get("number", Long.class)— получаем значение поляnumberиз входящего сообщения. В методMessage#getнеобходимо передавать класс значения для однозначного преобразования из сериализованной формы в Java-объект. -
messageBuilder.set("number_x2", number * 2)— записываем значение в полеnumber_x2. Это поле должно присутствовать в схеме стримаx2_numbersв статической спеке. -
output.addMessage(messageBuilder.finish())— методfinishвозвращает готовое сообщение, которое добавляется вOutputCollector.
BatchFunction
BatchFunction получает весь список сообщений и таймеров, пришедших от worker-а. Интерфейс предоставляет два метода:
onMessages(List<ExtendedMessage> messages, OutputCollector output, RuntimeContext ctx)— вызывается для батча сообщений.onTimers(List<Timer> timers, OutputCollector output, RuntimeContext ctx)— вызывается для батча таймеров.
Батч соответствует одному запросу воркера и может содержать сообщения с разными ключами; группировка по ключам при необходимости выполняется в пользовательском коде (см. Companion).
Пример batch-функции
public class X2BatchMapper implements BatchFunction {
@Override
public void onMessages(List<ExtendedMessage> messages, OutputCollector output, RuntimeContext ctx) {
var messageBuilder = ctx.createMessageBuilder("x2_numbers"); //1
for (var message : messages) { //2
Long number = message.get("number", Long.class); //3
messageBuilder.set("number_x2", number * 2); //4
output.addMessage(messageBuilder.finish()); //5
}
}
}
class X2BatchMapper : BatchFunction {
override fun onMessages(messages: List<ExtendedMessage>, output: OutputCollector, ctx: RuntimeContext) {
val messageBuilder = ctx.createMessageBuilder("x2_numbers") //1
for (message in messages) { //2
val number: Long? = message.get("number", Long::class.java) //3
messageBuilder.set("number_x2", number!! * 2) //4
output.addMessage(messageBuilder.finish()) //5
}
}
}
Ключевое отличие от RowFunction:
MessageBuilderсоздается один раз за весь батч (строка 1).- Метод
finish()одновременно с возвратом готового сообщения сбрасываетMessageBuilderв исходное состояние, после чего его можно переиспользовать для следующего сообщения (строка 5).
Регистрация в PipelineContext
Все объекты Computation и типизированные стримы (созданные через FlowStreams.typed) должны быть зарегистрированы в PipelineContext перед запуском GrpcServerExecution.
Нетипизированные стримы (созданные через FlowStreams.raw) регистрировать не обязательно, Flow сам создаст их на основе блока streams в статической спеке.
Подробнее про Типизированные стримы.
var context = new PipelineContext();
// Регистрация объектов Computation.
Computation join = Computation.builder()
.setComputationId("join")
.setProcessFunction(new JoinProcessFunction())
.build();
context.registerComputation(join);
SourceComputation reader = SourceComputation.builder()
.setComputationId("hit_reader")
.setProcessFunction(new HitParsingFunction())
.build();
context.registerComputation(reader);
// Регистрация типизированных стримов.
context.registerStream(FlowStreams.typed("hit", Hit.class));
context.registerStream(FlowStreams.typed("action", Action.class));
context.registerStream(FlowStreams.typed("joined_action", JoinedAction.class));
val context = PipelineContext()
// Регистрация объектов Computation.
val join: Computation = Computation.builder()
.setComputationId("join")
.setProcessFunction(JoinProcessFunction())
.build()
context.registerComputation(join)
val reader: SourceComputation = SourceComputation.builder()
.setComputationId("hit_reader")
.setProcessFunction(HitParsingFunction())
.build()
context.registerComputation(reader)
// Регистрация типизированных стримов.
context.registerStream(FlowStreams.typed("hit", Hit::class.java))
context.registerStream(FlowStreams.typed("action", Action::class.java))
context.registerStream(FlowStreams.typed("joined_action", JoinedAction::class.java))
Важно
Каждый Computation и стрим должен иметь уникальный идентификатор, соответствующий идентификаторам в статической спеке. Попытка зарегистрировать Computation или стрим с уже существующим идентификатором приведёт к ошибке и невозможности старта компаньона.
RuntimeContext
RuntimeContext предоставляет доступ к контексту выполнения компьютейшена. Основные методы:
| Метод | Описание |
|---|---|
ctx.createMessageBuilder(streamId) |
Создать MessageBuilder для указанного output-стрима |
ctx.getComputationParameters() |
Получить параметры компьютейшена из спеки |
ctx.getEpochInputEventWatermark() |
Получить текущий вотермарк эпохи |
ctx.getProtoStateAccessor(name, message, Class) |
Получить стейт в виде protobuf-объекта, привязанный к ключу сообщения |
ctx.getYsonStateAccessor(name, message, Class) |
Получить YSON-стейт, привязанный к ключу сообщения |
ctx.getStateAccessor(name, message, Class, ser, deser) |
Получить стейт с пользовательской сериализацией/десериализацией |
ctx.getRawStateAccessor(name, message) |
Получить стейт как массив байт без интерпретации |
ctx.getNoOpStateAccessor(name, message) |
Получить стейт, сохраняющий только факт наличия (без значения) |
ctx.getExternalStateAccessor(name, message) |
Получить внешний стейт, привязанный к ключу сообщения |
Подробнее про работу со стейтами — в разделе Работа со стейтами (Java).
OutputCollector
OutputCollector используется для отправки результатов обработки:
| Метод | Описание |
|---|---|
output.addMessage(message) |
Добавить выходное сообщение |
output.addMessage(message, options) |
Добавить сообщение с AddMessageOptions, задающим публикацию и суффикс message id для Swift |
output.addTimer(triggerTimestamp) |
Добавить таймер с указанным временем срабатывания (eventTimestamp = 0) |
output.addTimer(triggerTimestamp, eventTimestamp) |
Добавить таймер с указанным временем срабатывания и event-временем |
output.addTimer(timerStreamId, triggerTimestamp, eventTimestamp) |
Добавить таймер для конкретного timer-стрима |
output.setParentIds(parentIds) |
Задать parent ID для отслеживания lineage сообщений. Возвращает новый OutputCollector |
В Swift-компьютейшене MessageIdSuffix.payloadHash() или MessageIdSuffix.userDefined(value) позволяет отвязать идентичность сообщения от порядка порождения:
output.addMessage(
message,
AddMessageOptions.builder()
.setMessageIdSuffix(MessageIdSuffix.payloadHash())
.build());
Options по умолчанию используют порядковый номер сообщения. Хеш payload и пользовательский суффикс поддерживаются только Swift-компьютейшенами.
Spring Boot
При использовании Spring Boot компьютейшен регистрируется аннотацией @FlowComputation (или @FlowSourceComputation для источника) прямо на классе ProcessFunction. Аннотация мета-аннотирована @Component, поэтому класс автоматически становится Spring-бином:
@FlowComputation(id = "mapper")
public class WordCountMapper implements RowFunction {
@Override
public void onMessage(ExtendedMessage message, OutputCollector output, RuntimeContext ctx) {
// обработка сообщения
}
}
@FlowComputation(id = "mapper")
class WordCountMapper : RowFunction {
override fun onMessage(message: ExtendedMessage, output: OutputCollector, ctx: RuntimeContext) {
// обработка сообщения
}
}
Стримы объявляются как Spring-бины FlowStream<?> (либо через ComputationProvider.getStreams()):
@Configuration
public class StreamConfiguration {
@Bean
public FlowStream<Word> wordsStream() {
return FlowStreams.typed("words", Word.class);
}
}
@Configuration
class StreamConfiguration {
@Bean
fun wordsStream(): FlowStream<Word> = FlowStreams.typed("words", Word::class.java)
}
FlowStreams.typed(...) создаёт типизированный стрим, который автоматически сериализует и десериализует сообщения в Java-объекты указанного типа. Подробнее в разделе Typed Streams.
Подробнее про регистрацию через аннотации и ComputationProvider — в разделе Spring Boot интеграция.
Конфигурация ресурса CompanionManager
Для запуска компаньона на Java или Kotlin необходимо объявить ресурс CompanionManager в статической спеке:
"CompanionManager" = {
"resource_class_name" = "NYT::NFlow::NCompanion::TJavaCompanionManager";
"parameters" = {
"main_class" = "tech.ytsaurus.flow.examples.waitclickjoin.PipelineMain";
"timeout" = "10s";
"jdk_bin_path" = "/app/ytflow/jdk/bin/java";
"classpath" = "/app/ytflow/lib/*";
};
"dependencies" = {};
};
Параметр resource_class_name указывает на класс ресурса, который будет осуществлять запуск компаньона.
В случае компаньона на Java или Kotlin resource_class_name всегда должен быть NYT::NFlow::NCompanion::TJavaCompanionManager (поддерживает оба языка через JVM).
Подробнее про спеку в разделе Spec, DynamicSpec и Config.