Multiplexer в YTsaurus Flow (C++)
Multiplexer — шаблон процесс-функции, который по входному ключу читает связанный с ним набор записей и отдаёт каждую запись отдельным выходным сообщением. Типичный пример: на вход поступает ключ X, а на выходе нужны все строки сортированной динамической таблицы, у которых X является префиксом ключа.
Что обеспечивает базовый класс:
- равномерный прогресс между несколькими активными ключами (один большой ключ не блокирует обработку остальных);
- обработка коллапсов — если для уже обрабатываемого ключа прилетело новое входное сообщение, итерация перезапускается так, чтобы все строки были выданы с актуальной версией payload'а из входа.
Как устроена итерация
Базовый класс хранит состояние для каждого ключа с курсором (Offset) и при каждом срабатывании таймера вызывает у наследника FetchBatch, передавая текущий startOffsetExclusive. Наследник возвращает следующий курсор; если данные закончились — nullopt.
При коллапсе (повторное входное сообщение для активного ключа) базовый класс запоминает текущую позицию (InitialStartOffset) и читает данные в два прохода:
- дочитывает остаток от текущей позиции до конца (фаза 1);
- возвращается к началу и доходит до сохранённой точки коллапса (фаза 2).
Так гарантируется, что после коллапса будут выданы все строки с новой версией payload'а — включая те, которые уже были эмитены до коллапса со старой версией.
Готовый класс: TDynamicTableMultiplexerProcessFunction
Покрывает типовой сценарий: есть сортированная динтаблица, а на вход прилетают сообщения с ключом и payload'ом. На каждое входное сообщение нужно выдать по одной строке для каждой записи таблицы с этим ключом.
Параметры
{
"computation_class_name" = "NYT::NFlow::TProcessFunctionComputation";
"processing_function" = "TMyMultiplexerProcessFunction";
"processing_function_parameters" = {
"table_path" = "<cluster=primary>//path/to/lookup_table";
};
}
Зарегистрируйте наследника с теми же типами статических и динамических параметров:
YT_FLOW_DEFINE_PROCESS_FUNCTION(
TMyMultiplexerProcessFunction,
TDynamicTableMultiplexerParameters,
TDynamicMultiplexerParameters);
table_path — обязательный путь с указанием кластера. Перед первым чтением батча класс получает из схемы таблицы набор колонок для итерации (ключевые колонки после group_by_schema) и колонки данных, а затем кэширует их.
group_by_schema компьютейшена должна совпадать с ведущими ключевыми колонками таблицы. После этого префикса в таблице должна быть хотя бы одна дополнительная сортировочная ключевая колонка. Базовый класс сохраняет схему дополнительных ключевых колонок вместе со смещением и сбрасывает смещение, если эта схема изменилась между запусками пайплайна. Изменения только в колонках данных не приводят к такому сбросу. Схема таблицы кэшируется на всё время работы экземпляра процесс-функции.
В конфигурационном фрагменте выше показаны только поля процесс-функции. Для компьютейшена также нужны group_by_schema, входные и выходные потоки, поток таймеров типа current_time, зависимость таймера от потока таймеров и входного потока, а также allow_timer_self_dependency = %true. См. реализацию функции и её спецификацию пайплайна.
Что нужно реализовать
Унаследуйтесь от класса и откройте его конструктор через using TDynamicTableMultiplexerProcessFunction::TDynamicTableMultiplexerProcessFunction либо определите конструктор, который принимает TProcessFunctionContextPtr. Переопределите BuildOutputForRow, чтобы построить выходное сообщение для одной выбранной строки. Метод получает входной ключ, строку и её схему, пользовательское состояние, коллектор выходных сообщений и контекст выполнения. Если требуется передать данные из входного сообщения в выходное, задайте свой TUserState и переопределите OnInputMessage. Иначе оставьте TUserState по умолчанию (TEmptyMultiplexerUserState) и не переопределяйте OnInputMessage.
void BuildOutputForRow(
const TKey& key,
const TPayload& rowPayload,
const NTableClient::TTableSchemaPtr& rowSchema,
TStateAccessor<TUserState>& userState,
const IOutputCollectorPtr& output,
const IRuntimeContextPtr& context) override;
В полном примере по ссылкам выше используются вход (key, payload), таблица поиска [hash, key, secondary_key, region] и выход (key, secondary_key, region, payload).
rowPayload — одна строка таблицы целиком (без group_by-колонок) как TPayload, rowSchema описывает её колонки (имена + типы). Колонки достаются по имени через GetColumnValue<T> (payload.h).
Динамические параметры
Унаследованы от базового класса:
timer_period(по умолчанию 5 секунд) — как часто срабатывает таймер для ключа.batch_size(по умолчанию 1000) — размер одного батча, передаётся вLIMITзапроса.
В динамической спецификации компьютейшена укажите оба поля в processing_function_parameters:
{
"processing_function_parameters" = {
"timer_period" = "10s";
"batch_size" = 1000;
};
}
Базовый класс: TMultiplexerProcessFunction
Если источник данных — не сортированная динамическая таблица, а, например, структура в памяти или пользовательский RPC-сервис, наследуйтесь напрямую от TMultiplexerProcessFunction<TUserState> и реализуйте FetchBatch.
Контракт FetchBatch:
- Offset (
TKey) — должен быть сравним и монотонно возрастать в рамках одной итерации. Базовый класс проверит монотонность и упадёт, если контракт нарушен. startOffsetExclusive— читать данные строго после этой позиции.nulloptозначает «с самого начала».endOffsetInclusive— задаётся только в фазе 2 (после коллапса). Реализация должна не превышать эту границу — базовый класс это проверит.limit— рекомендуемый максимум строк в одном батче.- Возврат
nullopt— в текущем диапазоне данных больше нет. Базовый класс либо переключится в фазу 2, либо завершит итерацию.
OnInputMessage вызывается для каждого входного сообщения ключа — как при появлении нового ключа, так и при коллапсе. Если реализации нужно различать эти случаи, храните в TUserState отдельный признак инициализации. Не полагайтесь на userState.IsEmpty(): состояние, состоящее из значений по умолчанию, останется пустым. Если сохранять данные из входного сообщения не требуется, не переопределяйте этот метод.