---
metadata:
  - name: generator
    content: Diplodoc Platform v5.50.6
alternate:
  - https://ytsaurus.tech/docs/en/flow/connectors/queue.md
  - https://ytsaurus.tech/docs/ru/flow/connectors/queue.md
---
> **Documentation Index:** Fetch the complete configuration index at https://ytsaurus.tech/docs/ru/llms.txt

<!-- source: ru/_includes/flow/connectors/queue.md -->
# Коннектор QYT в YTsaurus Flow

Коннектор к [очередям YTsaurus](https://ytsaurus.tech/docs/ru/user-guide/dynamic-tables/queues.md) (QYT).

Код коннектора находится [здесь](https://github.com/ytsaurus/ytsaurus/tree/main/yt/yt/flow/library/cpp/connectors/queue).

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

При чтении из очереди простейшим способом в [`Computation`](https://ytsaurus.tech/docs/ru/flow/concepts/glossary.md#stream-and-computation) приходят сообщения той же схемы, что и у строк таблицы очереди.

По умолчанию `event_time` и `system_time` сообщения считаются на основе колонки `$timestamp`, которая содержит YT-таймстемп записи строки (подробнее в [документации упорядоченных динтаблиц](https://ytsaurus.tech/docs/ru/user-guide/dynamic-tables/queues.md)).
Но при наличии в очереди колонки с метой (смотри параметры статической спеки сорса), `event_time` сообщения может браться из соответствующего поля. Также из меты может браться информация о [ватермарке](https://ytsaurus.tech/docs/ru/flow/concepts/glossary.md#timestamps-and-watermarks) стрима.

Класс [сорса](https://ytsaurus.tech/docs/ru/flow/concepts/glossary.md#source): `NYT::NFlow::TQueueSource`.

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

{% note info %}

Источник определяется кластером и путём очереди (`queue_path`). При их изменении партиции пересоздаются и читаются заново — см. [Смена источника](https://ytsaurus.tech/docs/ru/flow/connectors/about.md#source-change).

{% endnote %}

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

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

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

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

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

Можно воспользоваться стандартной YT CLI, переставив значение оффсетов у консьюмера на желаемые для каждой партиции. Не очень удобно при первом чтении топика, так как необходимо знать значение самого свежего оффсета, и назначить его для каждой [партиции](https://ytsaurus.tech/docs/ru/flow/concepts/glossary.md#partition).

Пример команды:
```bash
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`.


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

<!-- source: ru/flow/generated_docs/NYT_NFlow_TUnitedParameters_NYT_NFlow_TQueueSource.md -->
<!-- This file is generated by yt/yt/flow/yandex/tools/generate_yson_struct_doc/generate.sh script -->
<!-- Before using doc generation tool check readme: yt/yt/flow/yandex/tools/generate_yson_struct_doc/README.md -->
Источник: [yt/yt/flow/library/cpp/common/registry-inl.h](https://github.com/ytsaurus/ytsaurus/tree/main/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>>>`
 ||
|#


{% cut "**Дополнительные параметры**" %}


#|
|| `update_info_period` | **Тип**: [TDuration](https://ytsaurus.tech/docs/ru/flow/generated_docs/all_yson_structs.md#TDuration)
**Значение по умолчанию**: `15s`
Период обновления служебной информации source о партиции. Конкретный коннектор может использовать этот тик для запросов статуса и проверки живости сессии. ||
|| `byte_size_alpha` | **Тип**: `double`
**Значение по умолчанию**: `0.05`
Коэффициент экспоненциального сглаживания средней суммы байтов и числа сообщений на один offset: чем он больше, тем быстрее оценка реагирует на новые данные. ||
|| `update_partition_count_period` | **Тип**: [TDuration](https://ytsaurus.tech/docs/ru/flow/generated_docs/all_yson_structs.md#TDuration)
**Значение по умолчанию**: `1m`
Как часто контроллеру обновлять число партиций в очереди (контроллер меняет число партиций computation в соответствии с числом партиций в очереди). ||
|| `update_partition_count_retry_min_backoff` | **Тип**: [TDuration](https://ytsaurus.tech/docs/ru/flow/generated_docs/all_yson_structs.md#TDuration)
**Значение по умолчанию**: `1s`
Начальная задержка перед повторной попыткой обновить число партиций после ошибки с учетом джиттера. Значение ограничивается `update_partition_count_period`; последующие задержки растут экспоненциально. ||
|#


{% endcut %}
<!-- endsource: ru/flow/generated_docs/NYT_NFlow_TUnitedParameters_NYT_NFlow_TQueueSource.md -->

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


<!-- source: ru/flow/generated_docs/NYT_NFlow_TDynamicUnitedParameters_NYT_NFlow_TQueueSource.md -->
<!-- This file is generated by yt/yt/flow/yandex/tools/generate_yson_struct_doc/generate.sh script -->
<!-- Before using doc generation tool check readme: yt/yt/flow/yandex/tools/generate_yson_struct_doc/README.md -->
Источник: [yt/yt/flow/library/cpp/common/registry-inl.h](https://github.com/ytsaurus/ytsaurus/tree/main/yt/yt/flow/library/cpp/common/registry-inl.h)

#|
|| **Параметр** | **Описание** ||
|| `pull_queue_timeout` | **Тип**: [TDuration](https://ytsaurus.tech/docs/ru/flow/generated_docs/all_yson_structs.md#TDuration)
**Значение по умолчанию**: `1m`
Таймаут для запроса чтения из очереди. ||
|#


{% cut "**Дополнительные параметры**" %}


#|
|| `unavailable_threshold` | **Тип**: [TDuration](https://ytsaurus.tech/docs/ru/flow/generated_docs/all_yson_structs.md#TDuration)
**Значение по умолчанию**: `5m`
Сколько времени подряд источник должен быть недоступен, чтобы партиция считалась стабильно недоступной. Засчитывается только то время, когда джоб работал и видел ошибку: простой между перезапусками в него не попадает, а любой успешный ответ источника обнуляет накопленное. ||
|#


{% endcut %}
<!-- endsource: ru/flow/generated_docs/NYT_NFlow_TDynamicUnitedParameters_NYT_NFlow_TQueueSource.md -->

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

Для записи сообщений в очередь в [синк](https://ytsaurus.tech/docs/ru/flow/concepts/glossary.md#sink) нужно отправлять сообщения той же схемы, что и у таблицы, в которую планируется запись. Также можно настроить запись метаинформации в специальную колонку, чтобы передавать информацию о `event_time` сообщений и о ватермарке стрима читателям очереди.

Есть два варианта синков: синхронный (`NYT::NFlow::TSyncQueueSink`) и асинхронный (`NYT::NFlow::TAsyncQueueSink`). Запись в синхронный синк идёт в основной транзакции [эпохи](https://ytsaurus.tech/docs/ru/flow/concepts/glossary.md#epoch), это эффективно, но можно писать только в очередь на основном кластере процессинга. Запись в асинхронный синк идёт уже после основной транзакции эпохи из сообщений сохранённых в output messages, это дороже, но так можно писать в очередь на любом кластере.

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

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

<!-- source: ru/flow/generated_docs/NYT_NFlow_TUnitedParameters_NYT_NFlow_TSyncQueueSink.md -->
<!-- This file is generated by yt/yt/flow/yandex/tools/generate_yson_struct_doc/generate.sh script -->
<!-- Before using doc generation tool check readme: yt/yt/flow/yandex/tools/generate_yson_struct_doc/README.md -->
Источник: [yt/yt/flow/library/cpp/common/registry-inl.h](https://github.com/ytsaurus/ytsaurus/tree/main/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`; если не задано ни одно из двух выражений &mdash; маршрутизация выключена (`$tablet_index` не пишется, таблет выбирает драйвер). Шардированные очереди неудобны в эксплуатации (жёсткая связность object&rarr;таблет, решардирование переразбивает ключи), поэтому включайте маршрутизацию только при реальной необходимости. Маршрутизация поддерживается только для синхронных 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](https://ytsaurus.tech/docs/ru/flow/generated_docs/all_yson_structs.md#NYT_NFlow_EQueueTabletIndexRoutingHashPolicy)`>`
Политика сведения хеша `tablet_index_routing_hash_expression` к индексу таблета. Обязателен вместе с ним. `range` &mdash; непрерывные равноширокие диапазоны хеша (rangeSize = 2^64 / tablet_count), рекомендуется: потребитель, партиционированный по тому же ключу диапазонами, читает только свой таблет. `modulo` &mdash; `hash % tablet_count`; не рекомендуется: `Computation`-ы Flow партиционированы по ключу диапазонами, поэтому очередь, шардированная по модулю, вынуждает каждого читателя читать все таблеты (full mesh на чтении). ||
|| `tablet_count` | **Тип**: `std::optional<long>`
Число таблетов для сведения `tablet_index_routing_hash_expression` к индексу таблета. Необязателен. Если не задан, синк периодически перечитывает `@tablet_count` целевой очереди (из кэша мастера) и подхватывает решардирование без рестарта. Если задан явно &mdash; значение фиксировано, запросов в очередь нет, а решардирование учитывается только при изменении конфигурации. ||
|| `column_filter` | **Тип**: `std::optional<THashSet<std::string>>`
Какие колонки из сообщения писать в очередь (по умолчанию &mdash; все). ||
|#


{% cut "**Дополнительные параметры**" %}


#|
|| `update_partition_count_period` | **Тип**: [TDuration](https://ytsaurus.tech/docs/ru/flow/generated_docs/all_yson_structs.md#TDuration)
**Значение по умолчанию**: `1m`
Как часто контроллеру обновлять число партиций в очереди (контроллер меняет число партиций computation в соответствии с числом партиций в очереди). ||
|| `update_partition_count_retry_min_backoff` | **Тип**: [TDuration](https://ytsaurus.tech/docs/ru/flow/generated_docs/all_yson_structs.md#TDuration)
**Значение по умолчанию**: `1s`
Начальная задержка перед повторной попыткой обновить число партиций после ошибки с учетом джиттера. Значение ограничивается `update_partition_count_period`; последующие задержки растут экспоненциально. ||
|#


{% endcut %}
<!-- endsource: ru/flow/generated_docs/NYT_NFlow_TUnitedParameters_NYT_NFlow_TSyncQueueSink.md -->

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

<!-- source: ru/flow/generated_docs/NYT_NFlow_TDynamicUnitedParameters_NYT_NFlow_TSyncQueueSink.md -->
<!-- This file is generated by yt/yt/flow/yandex/tools/generate_yson_struct_doc/generate.sh script -->
<!-- Before using doc generation tool check readme: yt/yt/flow/yandex/tools/generate_yson_struct_doc/README.md -->
Источник: [yt/yt/flow/library/cpp/common/registry-inl.h](https://github.com/ytsaurus/ytsaurus/tree/main/yt/yt/flow/library/cpp/common/registry-inl.h)

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


{% cut "**Дополнительные параметры**" %}


#|
|| `flow_queue_meta_heartbeat_period` | **Тип**: [TDuration](https://ytsaurus.tech/docs/ru/flow/generated_docs/all_yson_structs.md#TDuration)
**Значение по умолчанию**: `10s`
Как часто контроллер будет во все партиции очереди писать хартбиты с метой с ватермарками. ||
|#


{% endcut %}
<!-- endsource: ru/flow/generated_docs/NYT_NFlow_TDynamicUnitedParameters_NYT_NFlow_TSyncQueueSink.md -->

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

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

<!-- source: ru/flow/generated_docs/NYT_NFlow_TUnitedParameters_NYT_NFlow_TAsyncQueueSink.md -->
<!-- This file is generated by yt/yt/flow/yandex/tools/generate_yson_struct_doc/generate.sh script -->
<!-- Before using doc generation tool check readme: yt/yt/flow/yandex/tools/generate_yson_struct_doc/README.md -->
Источник: [yt/yt/flow/library/cpp/common/registry-inl.h](https://github.com/ytsaurus/ytsaurus/tree/main/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`; если не задано ни одно из двух выражений &mdash; маршрутизация выключена (`$tablet_index` не пишется, таблет выбирает драйвер). Шардированные очереди неудобны в эксплуатации (жёсткая связность object&rarr;таблет, решардирование переразбивает ключи), поэтому включайте маршрутизацию только при реальной необходимости. Маршрутизация поддерживается только для синхронных 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](https://ytsaurus.tech/docs/ru/flow/generated_docs/all_yson_structs.md#NYT_NFlow_EQueueTabletIndexRoutingHashPolicy)`>`
Политика сведения хеша `tablet_index_routing_hash_expression` к индексу таблета. Обязателен вместе с ним. `range` &mdash; непрерывные равноширокие диапазоны хеша (rangeSize = 2^64 / tablet_count), рекомендуется: потребитель, партиционированный по тому же ключу диапазонами, читает только свой таблет. `modulo` &mdash; `hash % tablet_count`; не рекомендуется: `Computation`-ы Flow партиционированы по ключу диапазонами, поэтому очередь, шардированная по модулю, вынуждает каждого читателя читать все таблеты (full mesh на чтении). ||
|| `tablet_count` | **Тип**: `std::optional<long>`
Число таблетов для сведения `tablet_index_routing_hash_expression` к индексу таблета. Необязателен. Если не задан, синк периодически перечитывает `@tablet_count` целевой очереди (из кэша мастера) и подхватывает решардирование без рестарта. Если задан явно &mdash; значение фиксировано, запросов в очередь нет, а решардирование учитывается только при изменении конфигурации. ||
|| `at_most_once_strategy` | **Тип**: `NYT::TIntrusivePtr<`[NYT::NFlow::TAtMostOnceStrategyParameters](https://ytsaurus.tech/docs/ru/flow/generated_docs/all_yson_structs.md#NYT_NFlow_TAtMostOnceStrategyParameters)`>`
Необязательная стратегия доставки at-most-once. Поддержка зависит от коннектора; перед включением проверьте документацию выбранного коннектора. ||
|| `column_filter` | **Тип**: `std::optional<THashSet<std::string>>`
Какие колонки из сообщения писать в очередь (по умолчанию &mdash; все). ||
|#


{% cut "**Дополнительные параметры**" %}


#|
|| `update_partition_count_period` | **Тип**: [TDuration](https://ytsaurus.tech/docs/ru/flow/generated_docs/all_yson_structs.md#TDuration)
**Значение по умолчанию**: `1m`
Как часто контроллеру обновлять число партиций в очереди (контроллер меняет число партиций computation в соответствии с числом партиций в очереди). ||
|| `update_partition_count_retry_min_backoff` | **Тип**: [TDuration](https://ytsaurus.tech/docs/ru/flow/generated_docs/all_yson_structs.md#TDuration)
**Значение по умолчанию**: `1s`
Начальная задержка перед повторной попыткой обновить число партиций после ошибки с учетом джиттера. Значение ограничивается `update_partition_count_period`; последующие задержки растут экспоненциально. ||
|#


{% endcut %}
<!-- endsource: ru/flow/generated_docs/NYT_NFlow_TUnitedParameters_NYT_NFlow_TAsyncQueueSink.md -->

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

<!-- source: ru/flow/generated_docs/NYT_NFlow_TDynamicUnitedParameters_NYT_NFlow_TAsyncQueueSink.md -->
<!-- This file is generated by yt/yt/flow/yandex/tools/generate_yson_struct_doc/generate.sh script -->
<!-- Before using doc generation tool check readme: yt/yt/flow/yandex/tools/generate_yson_struct_doc/README.md -->
Источник: [yt/yt/flow/library/cpp/common/registry-inl.h](https://github.com/ytsaurus/ytsaurus/tree/main/yt/yt/flow/library/cpp/common/registry-inl.h)

#|
|| **Параметр** | **Описание** ||
|| `at_most_once_strategy` | **Тип**: `NYT::TIntrusivePtr<`[NYT::NFlow::TAtMostOnceStrategyDynamicParameters](https://ytsaurus.tech/docs/ru/flow/generated_docs/all_yson_structs.md#NYT_NFlow_TAtMostOnceStrategyDynamicParameters)`>`
**Значение по умолчанию**: `{}`
Динамические параметры `at_most_once_strategy`. Поддержка зависит от коннектора; перед настройкой проверьте документацию выбранного коннектора. ||
|#


{% cut "**Дополнительные параметры**" %}


#|
|| `flow_queue_meta_heartbeat_period` | **Тип**: [TDuration](https://ytsaurus.tech/docs/ru/flow/generated_docs/all_yson_structs.md#TDuration)
**Значение по умолчанию**: `10s`
Как часто контроллер будет во все партиции очереди писать хартбиты с метой с ватермарками. ||
|| `write_period` | **Тип**: [TDuration](https://ytsaurus.tech/docs/ru/flow/generated_docs/all_yson_structs.md#TDuration)
**Значение по умолчанию**: `100ms`
 ||
|| `max_rows_per_write` | **Тип**: `long`
**Значение по умолчанию**: `1000`
 ||
|| `max_bytes_per_write` | **Тип**: `long`
**Значение по умолчанию**: `1048576`
 ||
|| `backoff_duration` | **Тип**: [TDuration](https://ytsaurus.tech/docs/ru/flow/generated_docs/all_yson_structs.md#TDuration)
**Значение по умолчанию**: `1s`
 ||
|#


{% endcut %}
<!-- endsource: ru/flow/generated_docs/NYT_NFlow_TDynamicUnitedParameters_NYT_NFlow_TAsyncQueueSink.md -->


## См. также

- [Список коннекторов](https://ytsaurus.tech/docs/ru/flow/connectors/about.md)
- [Spec и DynamicSpec](https://ytsaurus.tech/docs/ru/flow/concepts/spec.md)
<!-- endsource: ru/_includes/flow/connectors/queue.md -->
