Коннектор QYT в YTsaurus Flow

Коннектор к очередям YTsaurus (QYT).

Код коннектора находится здесь.

Чтение из очереди

При чтении из очереди простейшим способом в Computation приходят сообщения той же схемы, что и у строк таблицы очереди.

По умолчанию event_time и system_time сообщения считаются на основе колонки $timestamp, которая содержит YT-таймстемп записи строки (подробнее в документации упорядоченных динтаблиц).
Но при наличии в очереди колонки с метой (смотри параметры статической спеки сорса), event_time сообщения может браться из соответствующего поля. Также из меты может браться информация о ватермарке стрима.

Класс сорса: NYT::NFlow::TQueueSource.

Настройки сорса представлены ниже.

Примечание

Источник определяется кластером и путём очереди (queue_path). При их изменении партиции пересоздаются и читаются заново — см. Смена источника.

Управление лагом консьюмера

При регистрации консьюмера к очереди лаг для консьюмера устанавливается на последние сохранённые данные в очереди. Таким образом, после регистрации, если политика удаления старых строк на очереди равна 1 дню, лаг будет равняться также одному дню.

Такое поведение может привести к следующему сценарию: при выкатке релиза лаг данных будет равен дню, и обработка сообщений встанет, пока лаг не будет обработан.

Избежать этого можно несколькими способами:

1. Использование стандартной YT CLI

Можно воспользоваться стандартной YT CLI, переставив значение оффсетов у консьюмера на желаемые для каждой партиции. Не очень удобно при первом чтении топика, так как необходимо знать значение самого свежего оффсета, и назначить его для каждой партиции.

Пример команды:

yt advance-queue-consumer --proxy <cluster-name> //home/service/stable/consumer <cluster-name>://home/source_service/Data/queue --partition-index 0 --new-offset 486482013

Расширенные возможности доступны в документации YT и по параметру -h.

Статическая спека:

Источник: yt/yt/flow/library/cpp/common/registry-inl.h

Параметр

Описание

finite

Тип: bool
Значение по умолчанию: false
Считать source конечным, то есть запомнить число сообщений в нём при старте и перевести stream в состояние completed при вычитывании этого числа сообщений.

queue_path

Тип: NYT::NYPath::TRichYPath
Обязательный параметр
Путь до очереди с указанием кластера.

consumer_path

Тип: NYT::NYPath::TRichYPath
Обязательный параметр
Путь до консьюмера очереди с указанием кластера.

try_parse_flow_queue_meta

Тип: bool
Значение по умолчанию: true
Нужно ли брать и парсить мету из очереди (в мете лежат ватермарки предоставленные писателем в очередь).

flow_queue_meta_column

Тип: std::string
Значение по умолчанию: flow_queue_meta
Из какой колонки брать мету.

ignore_malformed_flow_queue_meta

Тип: bool
Значение по умолчанию: false
Игнорировать невалидную мету или падать.

partition_filter

Тип: std::optional<std::vector<std::pair<int, int>>>

Дополнительные параметры

update_info_period

Тип: TDuration
Значение по умолчанию: 15s
Период обновления служебной информации source о партиции. Конкретный коннектор может использовать этот тик для запросов статуса и проверки живости сессии.

byte_size_alpha

Тип: double
Значение по умолчанию: 0.05
Коэффициент экспоненциального сглаживания средней суммы байтов и числа сообщений на один offset: чем он больше, тем быстрее оценка реагирует на новые данные.

update_partition_count_period

Тип: TDuration
Значение по умолчанию: 1m
Как часто контроллеру обновлять число партиций в очереди (контроллер меняет число партиций computation в соответствии с числом партиций в очереди).

update_partition_count_retry_min_backoff

Тип: TDuration
Значение по умолчанию: 1s
Начальная задержка перед повторной попыткой обновить число партиций после ошибки с учетом джиттера. Значение ограничивается update_partition_count_period; последующие задержки растут экспоненциально.

Динамическая спека:

Источник: yt/yt/flow/library/cpp/common/registry-inl.h

Параметр

Описание

pull_queue_timeout

Тип: TDuration
Значение по умолчанию: 1m
Таймаут для запроса чтения из очереди.

Дополнительные параметры

unavailable_threshold

Тип: TDuration
Значение по умолчанию: 5m
Сколько времени подряд источник должен быть недоступен, чтобы партиция считалась стабильно недоступной. Засчитывается только то время, когда джоб работал и видел ошибку: простой между перезапусками в него не попадает, а любой успешный ответ источника обнуляет накопленное.

Запись в очередь

Для записи сообщений в очередь в синк нужно отправлять сообщения той же схемы, что и у таблицы, в которую планируется запись. Также можно настроить запись метаинформации в специальную колонку, чтобы передавать информацию о event_time сообщений и о ватермарке стрима читателям очереди.

Есть два варианта синков: синхронный (NYT::NFlow::TSyncQueueSink) и асинхронный (NYT::NFlow::TAsyncQueueSink). Запись в синхронный синк идёт в основной транзакции эпохи, это эффективно, но можно писать только в очередь на основном кластере процессинга. Запись в асинхронный синк идёт уже после основной транзакции эпохи из сообщений сохранённых в output messages, это дороже, но так можно писать в очередь на любом кластере.

Параметры спек синхронного синка

Статическая спека:

Источник: yt/yt/flow/library/cpp/common/registry-inl.h

Параметр

Описание

queue_path

Тип: NYT::NYPath::TRichYPath
Обязательный параметр
Путь до очереди с указанием кластера.

write_flow_queue_meta

Тип: bool
Значение по умолчанию: false
Нужно ли писать мету в очередь (в мете пишутся ватермарки).

flow_queue_meta_column

Тип: std::string
Значение по умолчанию: flow_queue_meta
В какую колонку писать мету.

tablet_index_expression

Тип: std::optional<std::string>
Verbatim-режим маршрутизации. QL-выражение (только встроенные функции, например farm_hash) над колонками сообщения, значение которого пишется в системную колонку $tablet_index как есть. Взаимоисключающе с tablet_index_routing_hash_expression; если не задано ни одно из двух выражений — маршрутизация выключена ($tablet_index не пишется, таблет выбирает драйвер). Шардированные очереди неудобны в эксплуатации (жёсткая связность object→таблет, решардирование переразбивает ключи), поэтому включайте маршрутизацию только при реальной необходимости. Маршрутизация поддерживается только для синхронных queue-синков; асинхронные синки отвергают параметры маршрутизации при загрузке спеки (YTFLOW-766).

tablet_index_routing_hash_expression

Тип: std::optional<std::string>
Hash-режим маршрутизации. QL-выражение (только встроенные функции, например farm_hash), вычисляющее uint64-хеш, который сводится к индексу таблета по tablet_index_routing_hash_policy над tablet_count. Взаимоисключающе с tablet_index_expression. Поддерживается только для синхронных queue-синков; асинхронные синки отвергают параметры маршрутизации (YTFLOW-766).

tablet_index_routing_hash_policy

Тип: std::optional<NYT::NFlow::EQueueTabletIndexRoutingHashPolicy>
Политика сведения хеша tablet_index_routing_hash_expression к индексу таблета. Обязателен вместе с ним. range — непрерывные равноширокие диапазоны хеша (rangeSize = 2^64 / tablet_count), рекомендуется: потребитель, партиционированный по тому же ключу диапазонами, читает только свой таблет. modulo — hash % tablet_count; не рекомендуется: Computation-ы Flow партиционированы по ключу диапазонами, поэтому очередь, шардированная по модулю, вынуждает каждого читателя читать все таблеты (full mesh на чтении).

tablet_count

Тип: std::optional<long>
Число таблетов для сведения tablet_index_routing_hash_expression к индексу таблета. Необязателен. Если не задан, синк периодически перечитывает @tablet_count целевой очереди (из кэша мастера) и подхватывает решардирование без рестарта. Если задан явно — значение фиксировано, запросов в очередь нет, а решардирование учитывается только при изменении конфигурации.

column_filter

Тип: std::optional<THashSet<std::string>>
Какие колонки из сообщения писать в очередь (по умолчанию — все).

Дополнительные параметры

update_partition_count_period

Тип: TDuration
Значение по умолчанию: 1m
Как часто контроллеру обновлять число партиций в очереди (контроллер меняет число партиций computation в соответствии с числом партиций в очереди).

update_partition_count_retry_min_backoff

Тип: TDuration
Значение по умолчанию: 1s
Начальная задержка перед повторной попыткой обновить число партиций после ошибки с учетом джиттера. Значение ограничивается update_partition_count_period; последующие задержки растут экспоненциально.

Динамическая спека:

Источник: yt/yt/flow/library/cpp/common/registry-inl.h

В структуре нет основных параметров.

Дополнительные параметры

flow_queue_meta_heartbeat_period

Тип: TDuration
Значение по умолчанию: 10s
Как часто контроллер будет во все партиции очереди писать хартбиты с метой с ватермарками.

Параметры спек асинхронного синка

Статическая спека:

Источник: yt/yt/flow/library/cpp/common/registry-inl.h

Параметр

Описание

queue_path

Тип: NYT::NYPath::TRichYPath
Обязательный параметр
Путь до очереди с указанием кластера.

write_flow_queue_meta

Тип: bool
Значение по умолчанию: false
Нужно ли писать мету в очередь (в мете пишутся ватермарки).

flow_queue_meta_column

Тип: std::string
Значение по умолчанию: flow_queue_meta
В какую колонку писать мету.

producer_path

Тип: NYT::NYPath::TRichYPath
Обязательный параметр
Путь до продюсера очереди с указанием кластера.

require_sync_replica

Тип: bool
Значение по умолчанию: true
Одноименный параметр при записи в очередь. Разрешена ли запись в очередь без синхронных реплик.

tablet_index_expression

Тип: std::optional<std::string>
Verbatim-режим маршрутизации. QL-выражение (только встроенные функции, например farm_hash) над колонками сообщения, значение которого пишется в системную колонку $tablet_index как есть. Взаимоисключающе с tablet_index_routing_hash_expression; если не задано ни одно из двух выражений — маршрутизация выключена ($tablet_index не пишется, таблет выбирает драйвер). Шардированные очереди неудобны в эксплуатации (жёсткая связность object→таблет, решардирование переразбивает ключи), поэтому включайте маршрутизацию только при реальной необходимости. Маршрутизация поддерживается только для синхронных queue-синков; асинхронные синки отвергают параметры маршрутизации при загрузке спеки (YTFLOW-766).

tablet_index_routing_hash_expression

Тип: std::optional<std::string>
Hash-режим маршрутизации. QL-выражение (только встроенные функции, например farm_hash), вычисляющее uint64-хеш, который сводится к индексу таблета по tablet_index_routing_hash_policy над tablet_count. Взаимоисключающе с tablet_index_expression. Поддерживается только для синхронных queue-синков; асинхронные синки отвергают параметры маршрутизации (YTFLOW-766).

tablet_index_routing_hash_policy

Тип: std::optional<NYT::NFlow::EQueueTabletIndexRoutingHashPolicy>
Политика сведения хеша tablet_index_routing_hash_expression к индексу таблета. Обязателен вместе с ним. range — непрерывные равноширокие диапазоны хеша (rangeSize = 2^64 / tablet_count), рекомендуется: потребитель, партиционированный по тому же ключу диапазонами, читает только свой таблет. modulo — hash % tablet_count; не рекомендуется: Computation-ы Flow партиционированы по ключу диапазонами, поэтому очередь, шардированная по модулю, вынуждает каждого читателя читать все таблеты (full mesh на чтении).

tablet_count

Тип: std::optional<long>
Число таблетов для сведения tablet_index_routing_hash_expression к индексу таблета. Необязателен. Если не задан, синк периодически перечитывает @tablet_count целевой очереди (из кэша мастера) и подхватывает решардирование без рестарта. Если задан явно — значение фиксировано, запросов в очередь нет, а решардирование учитывается только при изменении конфигурации.

at_most_once_strategy

Тип: NYT::TIntrusivePtr<NYT::NFlow::TAtMostOnceStrategyParameters>
Необязательная стратегия доставки at-most-once. Поддержка зависит от коннектора; перед включением проверьте документацию выбранного коннектора.

column_filter

Тип: std::optional<THashSet<std::string>>
Какие колонки из сообщения писать в очередь (по умолчанию — все).

Дополнительные параметры

update_partition_count_period

Тип: TDuration
Значение по умолчанию: 1m
Как часто контроллеру обновлять число партиций в очереди (контроллер меняет число партиций computation в соответствии с числом партиций в очереди).

update_partition_count_retry_min_backoff

Тип: TDuration
Значение по умолчанию: 1s
Начальная задержка перед повторной попыткой обновить число партиций после ошибки с учетом джиттера. Значение ограничивается update_partition_count_period; последующие задержки растут экспоненциально.

Динамическая спека:

Источник: yt/yt/flow/library/cpp/common/registry-inl.h

Параметр

Описание

at_most_once_strategy

Тип: NYT::TIntrusivePtr<NYT::NFlow::TAtMostOnceStrategyDynamicParameters>
Значение по умолчанию: {}
Динамические параметры at_most_once_strategy. Поддержка зависит от коннектора; перед настройкой проверьте документацию выбранного коннектора.

Дополнительные параметры

flow_queue_meta_heartbeat_period

Тип: TDuration
Значение по умолчанию: 10s
Как часто контроллер будет во все партиции очереди писать хартбиты с метой с ватермарками.

write_period

Тип: TDuration
Значение по умолчанию: 100ms

max_rows_per_write

Тип: long
Значение по умолчанию: 1000

max_bytes_per_write

Тип: long
Значение по умолчанию: 1048576

backoff_duration

Тип: TDuration
Значение по умолчанию: 1s

См. также

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