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
    }
}

Разберем построчно:

  1. ctx.createMessageBuilder("x2_numbers") — создается MessageBuilder для output-стрима с id = x2_numbers. Стрим с таким идентификатором должен присутствовать в списке output_stream_ids в статической спеке компьютейшена.

  2. message.get("number", Long.class) — получаем значение поля number из входящего сообщения. В метод Message#get необходимо передавать класс значения для однозначного преобразования из сериализованной формы в Java-объект.

  3. messageBuilder.set("number_x2", number * 2) — записываем значение в поле number_x2. Это поле должно присутствовать в схеме стрима x2_numbers в статической спеке.

  4. 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

Исходный код StatefulContext

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.

См. также

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