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

<!-- source: ru/_includes/flow/concepts/guarantees.md -->
# Гарантии обработки в YTsaurus Flow

## Обзор {#overview}

Flow по умолчанию обеспечивает **exactly-once семантику** обработки событий. Это означает:

- Потребители получают результат обработки каждого входного сообщения ровно один раз, без потерь и дубликатов.
- [Стейт](https://ytsaurus.tech/docs/ru/flow/concepts/stateful.md) обновляется атомарно вместе с обработкой сообщения.

При необходимости семантику можно ослабить до [**at-least-once**](#at-least-once) или [**at-most-once**](#at-most-once) — например, ради снижения задержек или потребления ресурсов пайплайном.

## Exactly-once семантика {#exactly-once}

### Как это работает {#how-it-works}

Exactly-once обеспечивается тремя механизмами, работающими совместно.

#### 1. Lease-транзакция {#lease}

`Lease` — мастер-транзакция, которой владеет контроллер; он создаёт её для каждой [джобы](https://ytsaurus.tech/docs/ru/flow/concepts/glossary.md#job). Пользователь не создаёт Lease и не управляет её временем жизни. Контроллер поддерживает транзакцию активной периодическими пингами и проверяет состояние всех Lease в каждом цикле планирования.

Каждая обычная транзакция, в которой джоба коммитит данные [эпохи](https://ytsaurus.tech/docs/ru/flow/concepts/glossary.md#epoch), указывает Lease этой джобы как пререквизит в `prerequisite_transaction_ids`. Для успешного коммита Lease должна оставаться активной. Lease завершается одним из двух способов. Контроллер отменяет её, когда удаляет джобу: при ребалансировке партиций, после сбоя джобы или потери воркера, а также при остановке, паузе или завершении пайплайна. Lease также может истечь сама, если контроллер прекращает поддерживать её пингами. В обоих случаях срабатывает один и тот же барьер.

После отмены или истечения Lease коммиты с этим пререквизитом отклоняются, поэтому зомби-джоба не может записать данные после того, как контроллер перестал считать её владельцем [партиции](https://ytsaurus.tech/docs/ru/flow/concepts/glossary.md#partition). Даже при отсутствии данных джоба периодически коммитит почти пустую транзакцию с тем же пререквизитом, чтобы обнаружить потерю Lease.

Если Lease истекла, контроллер замечает это в очередном цикле планирования и удаляет соответствующую джобу. Пока пайплайн остаётся активным и планирование продолжается, после отмены или истечения Lease контроллер удаляет старую джобу, если она ещё существует, планирует необходимые новые джобы для затронутых текущих партиций и создаёт для каждой новой джобы новую Lease. Новые джобы восстанавливают состояние из последней закоммиченной эпохи, возобновляют обработку затронутых текущих партиций и могут повторно обработать только незакоммиченную работу. При завершении, остановке или паузе пайплайна контроллер вместо этого удаляет джобы, отменяет их Lease и не планирует замену.

Lease служит барьером владения перед коммитом. Сама по себе она не выполняет [дедупликацию входных сообщений](#input-dedup), не обеспечивает [атомарность выходных данных](#output-atomicity) или [атомарность состояния](#automatic-guarantees) и не распространяется на [произвольные внешние побочные эффекты](#side-effects).

#### 2. Дедупликация входных сообщений {#input-dedup}

- **Внутренние потоки** (`input`): дедупликация по `message_id` с помощью внутренней таблицы `input_messages` или её компактного варианта `compact_input_messages`. Компактный вариант выбирается автоматически, когда не включён `experimental_enable_non_uint_key` (партиционирование идёт исключительно по `uint64`-хеш-колонке); поведение можно переопределить параметром `use_compact_input_messages`. Старые `message_id` удаляются из этой таблицы с использованием [SystemWatermark](https://ytsaurus.tech/docs/ru/flow/concepts/glossary.md#timestamps-and-watermarks).
- **Источники** (`source`): дедупликация по оффсетам — оффсеты консьюмера продвигаются только после успешного коммита [эпохи](https://ytsaurus.tech/docs/ru/flow/concepts/glossary.md#epoch).
- **[Таймеры](https://ytsaurus.tech/docs/ru/flow/concepts/glossary.md#timer)**: результат обработки таймера коммитится в той же транзакции, что и удаление этого таймера.

#### 3. Атомарность выходных данных {#output-atomicity}

Способ доставки выходных сообщений зависит от типа [синка](https://ytsaurus.tech/docs/ru/flow/concepts/glossary.md#sink):

- **Синхронные синки** (например, `TSyncQueueSink`) — записывают данные в целевую систему в той же транзакции эпохи. Запись и подтверждение атомарны.
- **Асинхронные синки** (например, `TAsyncQueueSink`) — синхронно сохраняют сообщения в YTsaurus в таблице `output_messages`, а отправляют сообщения в целевую систему уже после коммита эпохи, используя пару `(producer_id, seqNo)` для дедупликации на стороне приёмника (Queue API). Если отправка не удалась, сообщения из `output_messages` будут отправлены повторно с тем же `seqNo`. По сути реализуется паттерн [transactional outbox](https://en.wikipedia.org/wiki/Inbox_and_outbox_pattern#The_outbox_pattern) с дедупликацией в целевой системе, если она это поддерживает.

### Что требуется от разработчика {#developer-responsibilities}

Exactly-once гарантии Flow распространяются на внутреннее состояние пайплайна и встроенные механизмы доставки. Однако пользовательский код может нарушить эти гарантии, если не соблюдать следующие правила.

#### Детерминизм вычислений {#determinism}

Требования к детерминизму зависят от типа компьютейшена.

**Для [Swift](https://ytsaurus.tech/docs/ru/flow/concepts/swift.md)-компьютейшенов** (`TSwiftMapComputation`, `TSwiftOrderedSourceComputation`) детерминизм строго обязателен. В Swift-компьютейшенах результаты вычислений не материализуются между эпохами и могут быть пересчитаны при рестарте. Если разные попытки вычисления одного и того же сообщения дают разные результаты, возможно смешение результатов от разных попыток, потеря или дублирование данных.

**Для обычных компьютейшенов** (`TTransformComputation`, `TTransformOrderedSourceComputation`) недетерминированный код технически допустим: из всех попыток обработки (возникших из-за рестартов джоб) атомарно применяется ровно одна. Смешения результатов не происходит. Однако если код производит разные выходные данные при повторном запуске, это может привести к неожиданному поведению — особенно при наличии побочных эффектов (вроде записи во внешние системы) или при обновлении стейта на основе недетерминированных значений.

{% note warning "Что нарушает детерминизм" %}

- Использование текущего времени (`Now()`, `time.time()`) в бизнес-логике.
- Генерация случайных чисел.
- Обращение к внешним сервисам, результат которых может меняться между вызовами.

{% endnote %}

#### Побочные эффекты {#side-effects}

Любое действие за пределами Flow — HTTP-запрос, запись в файл, отправка уведомления — **не покрывается** exactly-once гарантиями. При рестарте джобы такое действие может быть выполнено повторно.

Если побочные эффекты неизбежны, разработчик должен самостоятельно обеспечить их идемпотентность (например, используя уникальный идентификатор сообщения как ключ идемпотентности).

#### Внешний стейт {#external-state}

Обновления [внутреннего стейта](https://ytsaurus.tech/docs/ru/flow/concepts/stateful.md) Flow (YSON-state, External State через `IExternalStateManager`) атомарны и безопасны. Однако если компьютейшен модифицирует данные в сторонней системе (базе данных, кеше), эти изменения не участвуют в транзакции эпохи и могут быть продублированы.

### Что Flow гарантирует автоматически {#automatic-guarantees}

- Внутренние потоки между [компьютейшенами](https://ytsaurus.tech/docs/ru/flow/concepts/glossary.md#computation) — exactly-once.
- Обновления [стейта](https://ytsaurus.tech/docs/ru/flow/concepts/stateful.md) — атомарны с обработкой сообщения.
- Обработка и удаление [таймера](https://ytsaurus.tech/docs/ru/flow/concepts/timers.md) — в одной транзакции.
- Продвижение оффсетов [сорсов](https://ytsaurus.tech/docs/ru/flow/concepts/glossary.md#source) — только после успешного коммита эпохи.

## At-least-once семантика {#at-least-once}

At-least-once гарантирует, что каждое сообщение будет обработано **хотя бы один раз**, но допускает повторную обработку. Это может быть полезно, когда бизнес-логика естественно идемпотентна, а стоимость дедупликации неоправданно высока.

### Режим `processing_mode: at_least_once_consistent` {#processing-mode}

Параметр `processing_mode` задаётся на уровне `TTransformComputation` в спеке. Значение по умолчанию — `exactly_once`.

```
"computations" = {
    "MyComputation" = {
        "processing_mode" = "at_least_once_consistent";
        ...
    };
};
```

**Что меняется**: отключается дедупликация входных сообщений по `message_id`. Сообщения могут быть обработаны повторно после рестарта джобы.

**Что не меняется**: стейт остаётся консистентным (отсюда слово «consistent» в названии режима). Транзакционность эпохи, Lease-механизм, атомарность обновления стейта — всё это продолжает работать.

**Когда использовать**: если накладные расходы на дедупликацию существенны, а бизнес-логика идемпотентна (например, `insert_or_assign` вместо `increment`).

### Батчинг в Swift-компьютейшенах: `allow_batching_with_relaxed_guarantees` {#swift-allow-batching-with-relaxed-guarantees}

Параметр `allow_batching_with_relaxed_guarantees` задаётся в блоке `parameters` компьютейшена типа `TSwiftMapComputation`. Значение по умолчанию — `%false`.

```
"computations" = {
    "MyBatcher" = {
        "parameters" = {
            "allow_batching_with_relaxed_guarantees" = %true;
        };
        ...
    };
};
```

**Что меняется**: одно выходное сообщение может склеивать несколько входных (батчинг) — например, чтобы свернуть множество мелких сообщений в одно крупное и снизить нагрузку по числу сообщений на партицию ниже по конвейеру. Склеенное сообщение получает `MessageId`, детерминированно выведенный из набора `MessageId` родителей, а входное сообщение считается обработанным только после доставки всех его детей.

**Что это означает для гарантий**: при рестарте границы склейки могут отличаться. Реплей, воспроизводящий склейку с тем же составом родителей, даёт тот же `MessageId` и дедуплицируется; но если родитель попал в склейку другого состава, её `MessageId` будет другим — нижестоящие компьютейшены должны быть готовы видеть содержимое каждого родителя более одного раза (at-least-once). Кроме того, порядок по `MessageId` в рамках одного ключа на стороне нижестоящих компьютейшенов не сохраняется: логика, полагающаяся на упорядоченность по `MessageId` внутри ключа, должна быть переписана с учётом этого.

**Когда использовать**: когда нужно укрупнить поток сообщений (батчинг/свёртка), а нижестоящая логика идемпотентна и не зависит от порядка `MessageId` внутри ключа. С выключенным флагом (по умолчанию) каждое выходное сообщение имеет ровно одного родителя и Swift-компьютейшен сохраняет exactly-once.


## At-most-once семантика {#at-most-once}

At-most-once гарантирует, что каждое сообщение будет обработано **не более одного раза**, но допускает потерю сообщений.

Эта семантика достигается настройкой `at_most_once_strategy` на асинхронных синках, которые её поддерживают.

```
"sinks" = {
    "my_sink" = {
        "at_most_once_strategy" = {
            "enabled" = %true;
            "total_queue_bytes_limit" = 104857600;  // 100 МБ
            "suspend_destruction_duration" = 60000;   // мс
        };
        ...
    };
};
```

**Как работает**: сообщения помещаются во внутреннюю очередь. Если очередь заполнена (`total_queue_bytes_limit`), новые сообщения **отбрасываются без ошибки**. Дедупликация и гарантии доставки отсутствуют.

**Когда использовать**: телеметрия, метрики, best-effort уведомления — сценарии, где потеря части сообщений допустима, а пропускная способность важнее надёжности.

## Гарантии порядка {#ordering}

Flow **не гарантирует** глобальный порядок обработки сообщений. Однако существуют гарантии для производных сообщений с совпадающими [ключами](https://ytsaurus.tech/docs/ru/flow/concepts/glossary.md#key) в рамках одной цепочки [lineage](https://ytsaurus.tech/docs/ru/flow/concepts/lineage.md).

Подробности — в разделе [Порядок обработки сообщений](https://ytsaurus.tech/docs/ru/flow/concepts/ordering.md).

## Влияние коннекторов на гарантии {#connectors}

Каждый [коннектор](https://ytsaurus.tech/docs/ru/flow/connectors/about.md) имеет свой набор гарантий, зависящий от типа подключения (source/sink) и варианта синка (sync/async/at-least-once).

### Queue (QYT) {#queue-guarantees}

- **Source**: exactly-once — оффсеты консьюмера продвигаются только после коммита эпохи.
- **Синхронный синк** (`TSyncQueueSink`): exactly-once — запись в очередь происходит в основной транзакции эпохи. Работает только на основном кластере процессинга.
- **Асинхронный синк** (`TAsyncQueueSink`): exactly-once — сообщения сохраняются в `output_messages`, затем доставляются с `producer_id` + `seqNo`. Queue API дедуплицирует повторы. Работает кросс-кластерно.

Подробнее — [Queue](https://ytsaurus.tech/docs/ru/flow/connectors/queue.md).

### Static Table {#static-table-guarantees}

- **Source**: exactly-once — дедупликация по диапазонам чтения.
- **Синк** (`TArrivalOrderTableSink`): exactly-once — таблица и прогресс коммитятся одной master-транзакцией; frontier партиции дедуплицирует replay, поэтому при частично покрытом replay записывается только непокрытый хвост без перезапуска job; callback доставки вызывается только после этого внешнего коммита и следующего коммита Flow.

Подробнее — [Static Table](https://ytsaurus.tech/docs/ru/flow/connectors/static-table.md).

### HTTP {#http-guarantees}

- **Source**: отсутствует.
- **Асинхронный синк** (`TAsyncHttpSink`): at-least-once — POST-запросы выполняются вне транзакции эпохи, поэтому получатель должен выполнять операцию идемпотентно или дедуплицировать стабильный идентификатор сообщения из заголовка `Idempotency-Key`, который используется по умолчанию. При включении `at_most_once_strategy` загрузка спеки завершается ошибкой.

Подробнее — [HTTP-расширение](https://ytsaurus.tech/docs/ru/flow/extensions/http.md).


### ClickHouse {#clickhouse-guarantees}

ClickHouse-расширение предоставляет три уровня гарантий в четырёх классах синков; выбор напрямую влияет на гарантии обработки сообщений:

- **Exactly-once синки** (`TClickHouseBatchingSink`, рекомендованный по умолчанию, и `TShardedClickHouseBatchingSink`): сообщения сохраняются в `output_messages` и детерминированно группируются в батчи по одним и тем же границам `MessageId`. В нешардированном синке `insert_deduplication_token` равен максимальному `MessageId` батча. При использовании `shard_hosts` каждый шард получает этот токен батча с суффиксом `:<имя шарда>`. Байт-идентичный реплей несёт **тот же** токен дедупликации даже после `group_by`-репартиционирования, поэтому exactly-once переживает репартиционирование. ClickHouse дедуплицирует повтор по токену.
- **At-least-once синк** (`TAtLeastOnceClickHouseSink`): батч эпохи пишется синхронно в ходе обработки без промежуточного хранения и без токена дедупликации. Ниже задержка и нагрузка на YTsaurus. При ошибке записи эпоха не коммитится, поэтому сообщения не теряются. Известные временные ошибки повторяются без ограничения числа попыток, а для неклассифицированных ошибок выполняется не более `max_insert_attempts` попыток вставки всего (по умолчанию 10). Постоянные ошибки сразу приводят к отказу. При сбое после успешной записи, но до коммита эпохи, тот же батч запишется повторно, поэтому возможны дубликаты.
- **At-most-once синк** (`TAtMostOnceClickHouseSink`): гарантия действует в обоих режимах стратегии. При статическом значении по умолчанию `at_most_once_strategy.enabled = false` ожидающие сообщения остаются в `output_messages` и доставляются по порядку, поэтому ошибку подключения или запуска сессии до начала `INSERT` можно повторить. Значение `true` включает независимую отправку каждого сообщения через ограниченную очередь в памяти без ожидания доставки. Её динамический лимит задаёт `at_most_once_strategy.total_queue_bytes_limit`; при отбрасывании сообщений из-за переполнения синк пишет предупреждение. В обоих режимах после начала `INSERT` синк подтверждает сообщение с ошибкой вместо повторной вставки. Каждое сообщение даёт отдельный `INSERT`, поэтому, чтобы не создавать множество мелких частей данных в ClickHouse, поток стоит укрупнять выше по пайплайну.

Exactly-once требует один из движков таблиц ClickHouse, дедуплицирующих вставляемые блоки (см. [Data Replication](https://clickhouse.com/docs/engines/table-engines/mergetree-family/replication) и [SharedMergeTree](https://clickhouse.com/docs/cloud/reference/shared-merge-tree) в документации ClickHouse):

- `ReplicatedMergeTree`
- `ReplicatedReplacingMergeTree`
- `ReplicatedSummingMergeTree`
- `ReplicatedAggregatingMergeTree`
- `SharedMergeTree`

Каждый экземпляр синка создаёт сессию записи отложенно при первой записи и сохраняет её до своей остановки. Полная проверка метаданных целевой таблицы выполняется при запуске этой сессии. При динамической перенастройке клиенты пересоздаются после изменения таймаута записи, а при изменении связанных параметров обновляется проверка окна дедупликации. Движок, схема и идентификатор репликации повторно не проверяются. После изменения целевой таблицы приостановите и снова запустите пайплайн, чтобы новые экземпляры синка проверили её. `Distributed`-таблицы, движки за пределами семейства `MergeTree` и цели с несколькими или резервными хостами внутри шарда, для которых нельзя доказать принадлежность одной логической таблице, отвергаются. Однохостовый нереплицируемый `MergeTree` без блочной дедупликации принимается с предупреждением, но exactly-once вырождается в at-least-once. Окно дедупликации конечно и вытесняется: если реплей доходит до ClickHouse позже вытеснения токена дедупликации, exactly-once также вырождается в at-least-once. Если известное окно дедупликации короче `replay_horizon`, синк пишет структурированное предупреждение с атрибутами `Database`, `Table`, `DedupWindowSetting`, `ReplayHorizon` и `DedupWindow`.

Сводная таблица гарантий по классам синка:

#|
|| | **Exactly-once синк** | **At-least-once синк** | **At-most-once синк** ||
|| Промежуточное хранение в YTsaurus | Да (`output_messages`) | Нет | Да при `at_most_once_strategy.enabled = false`; нет при `true` ||
|| Дедупликация в ClickHouse | `insert_deduplication_token` (max `MessageId`; с `shard_hosts` добавляется `:<имя шарда>`) | Нет | Нет ||
|| Задержка | Выше (доп. запись) | Ниже | Ниже ||
|| Потеря при сбое | Нет | Нет | Возможна после начала `INSERT`; при `enabled = true` также из-за переполнения очереди ||
|| Дубликаты при сбое | Нет | Возможны | Нет ||
|#

Подробнее — [ClickHouse-расширение](https://ytsaurus.tech/docs/ru/flow/extensions/clickhouse.md).

### Сервис-лог {#servicelog-guarantees}

- **Source**: exactly-once — дедупликация по диапазонам чтения. Сервис-лог — это бесконечный источник, циклически обходящий таблицу. Гарантии exactly-once применяются в рамках каждого цикла обхода.
- Синк отсутствует. Сортированная динамическая таблица, используемая как источник сервис-лога, зачастую наполняется самим пайплайном через работу с [внешним стейтом](https://ytsaurus.tech/docs/ru/flow/concepts/stateful.md#external-state). Например, компьютейшен типа `TTransformComputation` обновляет записи таблицы на каждой эпохе.

Подробнее — [Сервис-лог](https://ytsaurus.tech/docs/ru/flow/connectors/servicelog.md).


## Отказоустойчивость {#fault-tolerance}

- [Пайплайн](https://ytsaurus.tech/docs/ru/flow/concepts/glossary.md#pipeline) переживает выпадение отдельных машин и датацентров.
- Пока пайплайн остаётся активным и планирование продолжается, при сбое джобы контроллер удаляет её, отменяет [Lease](#lease), если та ещё активна, и планирует необходимые новые джобы для затронутых текущих партиций, создавая для каждой новой джобы новую Lease. Подробнее о барьере и восстановлении см. в разделе [Lease-транзакция](#lease).
- Внутренние потоки хранятся в динамических таблицах YTsaurus — потеря данных невозможна при штатной работе хранилища.
- [Автоматическая балансировка партиций](https://ytsaurus.tech/docs/ru/flow/about.md) перераспределяет нагрузку при изменении топологии кластера.

## См. также

- [Обзор Flow](https://ytsaurus.tech/docs/ru/flow/about.md)
- [Порядок обработки сообщений](https://ytsaurus.tech/docs/ru/flow/concepts/ordering.md)
- [Таймеры](https://ytsaurus.tech/docs/ru/flow/concepts/timers.md)
- [Stateful-обработка](https://ytsaurus.tech/docs/ru/flow/concepts/stateful.md)
- [Коннекторы](https://ytsaurus.tech/docs/ru/flow/connectors/about.md)
<!-- endsource: ru/_includes/flow/concepts/guarantees.md -->
