Коннектор к очередям 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
|
См. также