---
metadata:
  - name: generator
    content: Diplodoc Platform v5.50.4
alternate:
  - https://ydb.tech/docs/en/dev/streaming-query/guarantees.md
  - https://ydb.tech/docs/ru/dev/streaming-query/guarantees.md
sourcePath: ru/core/dev/streaming-query/guarantees.md
---
> **Documentation Index:** Fetch the complete configuration index at https://ydb.tech/docs/ru/llms.txt

# Гарантии доставки данных

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

{% note info %}

Мы постоянно работаем над развитием механизмов потоковой обработки. В будущих версиях предоставляемые гарантии будут улучшены.

{% endnote %}

**Гарантии обработки данных (dataplane)**:

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

**Аномалии при модификации запросов (control plane)**:

- [Потеря событий при пересоздании запроса](#incomplete-windows-restart) — при пересоздании запроса через DROP + CREATE часть событий, поступивших между удалением и созданием, будет пропущена.
- [Частичное первое окно агрегаций](#partial-first-window) — при старте запроса первое окно агрегации содержит неполные данные.

## Чекпоинты и восстановление {#checkpoints}

YDB периодически сохраняет [чекпоинт](https://ydb.tech/docs/ru/dev/streaming-query/checkpoints.md) — снимок состояния запроса, содержащий:

- [смещения](https://ydb.tech/docs/ru/concepts/datamodel/topic.md#consumer-offset) во входных топиках — позиции, до которых события были прочитаны и обработаны;
- состояния агрегаций — промежуточные результаты операций, например накопленные значения в [GROUP BY HOP](https://ydb.tech/docs/ru/yql/reference/syntax/select/group-by.md#group-by-hop).

YDB хранит смещения чтения в собственных чекпоинтах, а не полагается на смещения [потребителя (consumer)](https://ydb.tech/docs/ru/concepts/datamodel/topic.md#consumer) во внешней системе.

При восстановлении запрос откатывается к последнему чекпоинту: возобновляет чтение с сохранённых смещений и восстанавливает состояния агрегаций. События, поступившие между чекпоинтом и сбоем, будут обработаны повторно. Подробнее о механизме чекпоинтов — в разделе [Чекпоинты](https://ydb.tech/docs/ru/dev/streaming-query/checkpoints.md).

## Гарантии обработки данных (dataplane) — at-least-once {#at-least-once}

Если в процессе обработки потока происходит сбой (перезапуск вычислительного узла, сетевой разрыв, таймаут), YDB автоматически восстанавливает запрос из последнего чекпоинта. Гарантия [at-least-once](https://en.wikipedia.org/wiki/Reliable_messaging#At-least-once_delivery) обеспечивается для всех типов потоковых запросов — каждое событие будет обработано как минимум один раз. Запрос возобновляет чтение с сохранённого смещения и отправляет результаты обработки повторно. Это относится ко всем видам запросов: к запросам без агрегации (фильтрация, обогащение, трансформация), и к запросам с [оконной агрегацией](https://ydb.tech/docs/ru/yql/reference/syntax/select/group-by.md#group-by-hop).


```mermaid
sequenceDiagram
    participant Топик
    participant Запрос as Запрос<br/>GROUP BY HOP (1 мин)
    participant Приемник

    Note over Запрос: Чекпоинт сохранён<br/>смещение = 2, sum = 10
    Топик->>Запрос: value = 3 (смещение 3)
    Note over Запрос: sum = 13
    Топик->>Запрос: value = 7 (смещение 4)
    Note over Запрос: sum = 20
    Запрос-xЗапрос: Сбой обработки
    Note over Запрос: Восстановление из чекпоинта<br/>смещение = 2, sum = 10
    Топик->>Запрос: value = 3 (повторно)
    Note over Запрос: sum = 13
    Топик->>Запрос: value = 7 (повторно)
    Note over Запрос: sum = 20
    Note over Запрос: Окно закрыто
    Запрос->>Приемник: sum = 20
```

При записи результата в таблицу через [UPSERT](https://ydb.tech/docs/ru/yql/reference/syntax/upsert_into.md) повторная обработка не приводит к дублированию: UPSERT обновляет существующую строку по первичному ключу. Данные не теряются, дубли не накапливаются.

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

## Гарантии при модификации запроса (control plane) {#modification-anomalies}

В настоящий момент изменение текста запроса без его остановки не поддерживается. Для обновления запроса используется сочетание команд [DROP](https://ydb.tech/docs/ru/yql/reference/syntax/drop-streaming-query.md) + [CREATE](https://ydb.tech/docs/ru/yql/reference/syntax/create-streaming-query.md), в этом случае гарантия `at-least-once` не выполняется: часть событий может быть пропущена. Ниже описаны сценарии, в которых это происходит.

### Частичные результаты первого окна при старте запроса {#partial-first-window}

Временные окна ([GROUP BY HOP](https://ydb.tech/docs/ru/yql/reference/syntax/select/group-by.md#group-by-hop)) рассчитывают свои границы по абсолютному (wall-clock) времени. Границы окон выравниваются по кратным интервалам от начала эпохи: например, при окне в 1 минуту границы всегда проходят в 12:00:00, 12:01:00, 12:02:00 и т.д., независимо от того, когда запрос был запущен. Если запрос стартует в 12:00:30, он попадает в уже идущее окно [12:00:00 .. 12:01:00], но данные начинают поступать только с 12:00:30. В результате агрегат первого окна вычисляется по данным за 30 секунд вместо полной минуты.

```mermaid
sequenceDiagram
    participant Топик
    participant Запрос as Запрос<br/>GROUP BY HOP (1 мин)
    participant Приемник

    Note over Запрос: Запрос стартует в 12:00:30
    Note over Топик: Первое событие приходит в 12:00:35
    Note over Запрос: Открывается окно [12:00:00 .. 12:01:00]
    Топик->>Запрос: value = 5 (12:00:35)
    Топик->>Запрос: value = 3 (12:00:42)
    Топик->>Запрос: value = 8 (12:00:55)
    Note over Запрос: Окно закрыто в 12:01:00
    Запрос->>Приемник: sum = 16 (за 25 сек вместо 60 сек)
```

Это ожидаемое поведение при первом запуске — все последующие окна получат данные за полный интервал, который важно учитывать при пересоздании запроса.

### Потеря событий при пересоздании запроса {#incomplete-windows-restart}

Для изменения текста запроса используется сочетание команд [DROP](https://ydb.tech/docs/ru/yql/reference/syntax/drop-streaming-query.md) + [CREATE](https://ydb.tech/docs/ru/yql/reference/syntax/create-streaming-query.md). При `DROP` чекпоинт удаляется вместе с запросом, так как YDB использует внутреннее хранение смещений чтения из источника, то эти смещения удаляются вместе с запросом. Новый запрос не имеет сохранённой позиции и начинает чтение с конца топика. Все события, поступившие в топик между удалением старого запроса и стартом нового, не будут прочитаны.

Аналогичная ситуация возникает, если данные, на которые указывает смещение в чекпоинте, уже удалены из топика по [TTL](https://ydb.tech/docs/ru/concepts/datamodel/topic.md#retention-time).

```mermaid
sequenceDiagram
    participant Топик
    participant Запрос v1
    participant Запрос v2

    Топик->>Запрос v1: События A..D
    Note over Запрос v1: Чекпоинт: смещение = 4
    Note over Запрос v1: DROP STREAMING QUERY<br/>(чекпоинт удалён)
    Note over Топик: События E, F поступают в топик
    Note over Запрос v2: CREATE STREAMING QUERY<br/>(старт с конца топика)
    Топик--xЗапрос v2: E, F (не прочитаны)
    Топик->>Запрос v2: G
    Топик->>Запрос v2: H
    Note over Запрос v2: Агрегат окна = SUM(G, H)<br/>вместо SUM(E..H)
```

Для запросов с оконной агрегацией первые окна после пересоздания будут содержать пропуски данных и заниженные агрегаты.

## См. также

- [Потоковые запросы](https://ydb.tech/docs/ru/concepts/streaming-query/streaming-query.md) — общее описание потоковых запросов.
- [Чекпоинты](https://ydb.tech/docs/ru/dev/streaming-query/checkpoints.md) — механизм чекпоинтов, обеспечивающий восстановление после сбоев.
- [Запись в таблицы](https://ydb.tech/docs/ru/dev/streaming-query/table-writing.md) — запись в таблицы и идемпотентность UPSERT.
