Computation в YTsaurus Flow

Computation — основной строительный блок пайплайна. Каждый Computation получает сообщения из входных потоков, обрабатывает их и отправляет результаты в выходные потоки.

Виды Computation

Во Flow реализовано четыре базовых вида Computation, каждый из которых описан в следующих разделах.

Классы с Swift в названии реализуют принцип Swift — подход к обработке данных без полной материализации, но с сохранением exactly-once гарантий и требованием к детерминированности преобразований.

В C++ эти классы задают режим исполнения и используются встроенными адаптерами. Новую пользовательскую логику реализуйте только как process function, а в спеке выбирайте подходящий TProcessFunction*Computation. Не создавайте пользовательские классы-наследники от базовых классов Computation.

TTransformComputation

Режим для произвольных преобразований входных данных. Результат обработки сохраняется в YTsaurus, поэтому нет требований к детерминированности. Поддерживает таймеры, стейты и Sink. Пользовательскую process function запускает в этом режиме TProcessFunctionComputation; для passthrough-варианта без бизнес-логики используют TPassthroughComputation.

TSwiftMapComputation

Режим детерминированного Map без материализации результатов в YTsaurus. Не поддерживает таймеры, Source и Sink. Process function должна быть строго детерминированной — при необходимости результат будет вычислен повторно. Для пользовательской логики используют TProcessFunctionSwiftMapComputation, для passthrough-варианта — TSwiftPassthroughComputation.

TSwiftOrderedSourceComputation

Source-режим для чтения данных из внешних источников. Требует, чтобы поток данных из каждого инстанса был упорядочен. Поддерживает WatermarkStrategy для оценки вотермарков. Для пользовательской логики используют TProcessFunctionSourceComputation, для passthrough-варианта — TSwiftPassthroughOrderedSourceComputation.

TTransformOrderedSourceComputation

Режим для обработки данных Source произвольной пользовательской логикой: парсинга, фильтрации или разворачивания одного сообщения в несколько. Пользовательскую process function запускает в этом режиме TProcessFunctionTransformOrderedSourceComputation, заменяя связку TSwiftPassthroughOrderedSourceComputation → TProcessFunctionComputation.

Результат обработки материализуется в YTsaurus, как у TTransformComputation, поэтому требований к детерминированности нет: после рестарта Flow доставляет уже материализованные сообщения с ранее назначенными им MessageId, а не вычисляет их заново. Смещение источника, материализованные выходные сообщения и стейты коммитятся в одной транзакции YTsaurus — обработка каждого сообщения источника применяется ровно один раз, включая обновления стейта.

Собственный стейт process function хранит в поле TMutableStateKeyClient<T>, инициализирует через initContext->InitClient(...) в Init(const IRuntimeInitContextPtr&) и читает через GetState(message->Key) при обработке. Пример — в разделе Process function (C++).

Поддерживаются source_streams (ровно один упорядоченный Source), несколько выходных стримов, watermark_strategy (watermark_generator оценивает вотермарки источника, watermark_alignment выравнивает чтение, event_timestamp_assigner назначает event_timestamp), skip_if_expression и сообщения с distribute = false. Непустой group_by_schema, input-стримы, таймеры и key-visitor-стримы приводят к ошибке валидации спеки.

Passthrough Computation

Passthrough-компьютейшен не содержит пользовательской бизнес-логики: входящие сообщения конвертируются в схему выходного стрима и передаются дальше без изменений. Используется для простого приведения схем между стримами, например при чтении очереди и перекладывании данных в другой стрим без какой-либо обработки.

В Flow реализованы три C++-класса:

Класс Базовый класс Назначение
TPassthroughComputation TTransformComputation Конвертирует input-сообщения в схему output-стрима
TSwiftPassthroughComputation TSwiftMapComputation Аналогично, без материализации (Swift)
TSwiftPassthroughOrderedSourceComputation TSwiftOrderedSourceComputation Конвертирует source-сообщения в output-стрим

Passthrough реализуется в Flow нативно на C++ и не требует Java- или Python-компаньона. Чтобы включить его, в статической спеке компьютейшена укажите соответствующий C++-класс в поле computation_class_name:

"passthrough" = {
    "computation_class_name" = "NYT::NFlow::TPassthroughComputation";
    "group_by_schema" = [...];
    "input_stream_ids" = [...];
    "output_stream_ids" = [...];
};

Подробнее — Computation (C++).

Общие свойства

  • Всё выполнение в рамках одной партиции строго однопоточно. Многопоточность достигается за счёт увеличения числа партиций.
  • Все Computation берут на себя заполнение метаполей message и timer.
  • Объект OutputCollector предназначен для сбора выходных сообщений и таймеров.
  • Метод SetParents позволяет управлять lineage сообщений для корректного расчёта метаполей.

Реализация на разных языках

Каждый язык предоставляет свой набор интерфейсов для реализации Computation:

  • C++: реализация process function (IProcessFunction, IBatchProcessFunction или IKeyedBatchProcessFunction) и выбор встроенного адаптера режима в спеке. Подробнее →
  • Java: реализация интерфейсов RowFunction или BatchFunction с методами onMessage/onTimer. Подробнее →
  • Python: наследование от RowFunction или BatchFunction с методами on_message/on_timer. Подробнее →
  • Go: реализация интерфейсов flow.RowFunction (OnMessage) или flow.BatchFunction (OnMessages); таймеры — отдельными интерфейсами flow.RowTimerFunction/flow.BatchTimerFunction. Подробнее →
  • YQL: компьютейшны генерируются автоматически по декларативному описанию. Подробнее →

См. также

Фильтрация входа

skip_if_expression фильтрует входные сообщения до дедупликации и пользовательской обработки. Пропущенные сообщения не записываются в хранилище дедупликации. Отфильтрованные сообщения подтверждаются без вызова обработчика, в том числе если отфильтрован весь батч. Статистика входа, включая распределение ключей и heavy hitters, описывает исходный поток до фильтрации.

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