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: компьютейшны генерируются автоматически по декларативному описанию. Подробнее →
См. также
- Stateful processing
- Вотермарки
- Таймеры
- Спеки
- Process function (C++)
- Режимы Computation (C++)
- Computation (Java)
- Computation (Python)
- Computation (Go)
- Computation (YQL)
Фильтрация входа
skip_if_expression фильтрует входные сообщения до дедупликации и пользовательской обработки. Пропущенные сообщения не записываются в хранилище дедупликации. Отфильтрованные сообщения подтверждаются без вызова обработчика, в том числе если отфильтрован весь батч. Статистика входа, включая распределение ключей и heavy hitters, описывает исходный поток до фильтрации.