Multiplexer в YTsaurus Flow (C++)

Multiplexer — шаблон процесс-функции, который по входному ключу читает связанный с ним набор записей и отдаёт каждую запись отдельным выходным сообщением. Типичный пример: на вход поступает ключ X, а на выходе нужны все строки сортированной динамической таблицы, у которых X является префиксом ключа.

Что обеспечивает базовый класс:

  • равномерный прогресс между несколькими активными ключами (один большой ключ не блокирует обработку остальных);
  • обработка коллапсов — если для уже обрабатываемого ключа прилетело новое входное сообщение, итерация перезапускается так, чтобы все строки были выданы с актуальной версией payload'а из входа.

Как устроена итерация

Базовый класс хранит состояние для каждого ключа с курсором (Offset) и при каждом срабатывании таймера вызывает у наследника FetchBatch, передавая текущий startOffsetExclusive. Наследник возвращает следующий курсор; если данные закончились — nullopt.

При коллапсе (повторное входное сообщение для активного ключа) базовый класс запоминает текущую позицию (InitialStartOffset) и читает данные в два прохода:

  1. дочитывает остаток от текущей позиции до конца (фаза 1);
  2. возвращается к началу и доходит до сохранённой точки коллапса (фаза 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(): состояние, состоящее из значений по умолчанию, останется пустым. Если сохранять данные из входного сообщения не требуется, не переопределяйте этот метод.

См. также

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