---
metadata:
  - name: generator
    content: Diplodoc Platform v5.52.0
alternate:
  - https://ydb.tech/docs/en/concepts/streaming-query/streaming-query.md?version=v26.1
  - https://ydb.tech/docs/ru/concepts/streaming-query/streaming-query.md?version=v26.1
  - href: ru/concepts/streaming-query/streaming-query.md
    type: text/markdown
    title: Markdown version
  - href: ../../llms.txt
    type: text/markdown
    title: llms.txt
sourcePath: ru/core/concepts/streaming-query/streaming-query.md
---
> **Documentation Index:** Fetch the complete configuration index at https://ydb.tech/docs/ru/llms.txt

# Потоковые запросы

Потоковый запрос — это тип запроса, предназначенный для непрерывной обработки неограниченного потока данных ([stream processing](https://en.wikipedia.org/wiki/Stream_processing)). В отличие от обычных запросов, потоковый запрос не завершается после получения результата, а работает постоянно, обрабатывая данные по мере их поступления.

Потоковая обработка данных — широко применяемый подход, реализованный в таких системах, как [Apache Flink](https://flink.apache.org/), [Apache Kafka Streams](https://kafka.apache.org/documentation/streams/), [Amazon Kinesis Data Analytics](https://aws.amazon.com/kinesis/data-analytics/). Типичные сценарии применения — мониторинг и алертинг, агрегация метрик по временным окнам, трансформация и фильтрация событий на лету, обогащение событий данными из справочников, поиск паттернов в последовательностях событий.

YDB реализует потоковую обработку как часть единой платформы для работы с данными. В дополнение к типовым сценариям использования потоковых запросов, интеграция в общую платформу YDB позволяет получать данные из [топиков](https://ydb.tech/docs/ru/concepts/datamodel/topic.md?version=v26.1) YDB, обрабатывать потоки изменений из [строковых таблиц](https://ydb.tech/docs/ru/concepts/datamodel/table.md?version=v26.1#row-oriented-tables) через [CDC](https://ydb.tech/docs/ru/concepts/cdc.md?version=v26.1), а результаты обработки записывать в выходные топики или напрямую в таблицы. Это позволяет строить конвейеры обработки данных внутри YDB с минимальными временными задержками.

## Отличия от обычных запросов {#differences}

Обычные запросы работают с данными, которые уже сохранены в таблицах. Запрос выполняется, возвращает результат и завершается. Потоковый запрос создается и продолжает работу бесконечно до явной отмены пользователем. Данные непрерывно поступают в топик, «протекают» через запрос и записываются в приёмник — другой топик или таблицу.


| Характеристика | Обычные запросы | Потоковые запросы |
|----------------|-----------------|-------------------|
| Данные | Конечный набор в таблицах | Бесконечный поток событий |
| Время жизни | Завершается после обработки | Работает непрерывно |
| Результат | Доступен после завершения | Обновляется по мере поступления данных |
| Восстановление при сбоях | Перезапуск вручную | Автоматическое восстановление из [чекпоинта](../../dev/streaming-query/checkpoints.md) |

## Источники и приёмники данных {#data-flow}

Потоковые запросы читают данные из [топиков](https://ydb.tech/docs/ru/concepts/datamodel/topic.md?version=v26.1) YDB и записывают результаты в [топики](https://ydb.tech/docs/ru/concepts/datamodel/topic.md?version=v26.1) или [таблицы](https://ydb.tech/docs/ru/concepts/datamodel/table.md?version=v26.1) YDB. Данные в топик могут поступать из внешних систем (например, приложение записывает события через SDK), но сам потоковый запрос работает только с сущностями YDB. Прямое чтение из внешних систем, например из Apache Kafka, или запись в таблицы других баз данных, таких как PostgreSQL, не поддерживается.

### Источники {#sources}

**Топики** — основной источник потоковых данных. Запрос читает сообщения из одного или нескольких [топиков](https://ydb.tech/docs/ru/concepts/datamodel/topic.md?version=v26.1) и обрабатывает их по мере поступления. Данные в топик могут записывать:

* **Внешние приложения** — через [YDB SDK](https://ydb.tech/docs/ru/reference/ydb-sdk/index.md?version=v26.1) или [Kafka API](https://ydb.tech/docs/ru/reference/kafka-api/index.md?version=v26.1). Например, сервис отправляет события телеметрии или логи в топик YDB, а потоковый запрос обрабатывает их.
* **CDC (Change Data Capture)** — потоки изменений из таблиц, реализуемые через встроенные [топики](https://ydb.tech/docs/ru/concepts/datamodel/topic.md?version=v26.1). Позволяют реагировать на вставки, обновления и удаления записей в реальном времени. Подробнее: [Change Data Capture (CDC)](https://ydb.tech/docs/ru/concepts/cdc.md?version=v26.1).

### Приёмники {#sinks}

**Топики** — для передачи результатов другим системам или следующим этапам обработки.

**Таблицы** — для материализации результатов. Данные сохраняются через [UPSERT](https://ydb.tech/docs/ru/yql/reference/syntax/upsert_into.md?version=v26.1) и доступны для обычных SQL-запросов.

## Гарантии {#guarantees}

В штатном режиме работы потоковые запросы обеспечивают гарантию доставки данных [at-least-once](https://en.wikipedia.org/wiki/Reliable_messaging#At-least-once_delivery) — при сбоях запрос автоматически восстанавливается из [чекпоинта](https://ydb.tech/docs/ru/dev/streaming-query/checkpoints.md?version=v26.1) и возобновляет чтение с сохранённых смещений.

При этом важно учитывать ограничения текущей реализации:

- **Сброс смещений при пересоздании запроса:** изменение текста запроса через [DROP](https://ydb.tech/docs/ru/yql/reference/syntax/drop-streaming-query.md?version=v26.1) + [CREATE](https://ydb.tech/docs/ru/yql/reference/syntax/create-streaming-query.md?version=v26.1) приводит к удалению чекпоинта. Новый запрос начинает чтение с конца топика, пропуская события, поступившие в интервале между удалением старой версии и стартом новой.
- **Неполнота агрегатов из-за отсутствия watermarks:** из-за отсутствия механизма [watermarks](https://en.wikipedia.org/wiki/Watermark_(data_synchronization)) временные окна закрываются по системному времени (wall-clock). События, поступившие с задержкой (например, из-за медленных партиций или сетевых лагов), не попадают в агрегаты уже закрытых окон.

Подробное описание гарантий, аномалий и способов их минимизации — в разделе [Гарантии доставки данных](https://ydb.tech/docs/ru/dev/streaming-query/guarantees.md?version=v26.1).

{% note info %}

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

{% endnote %}

## Ограничения {#limitations}

{% note warning %}

- Запрос должен содержать хотя бы одно чтение из топика, так как потоковая обработка требует наличия непрерывного входного потока данных.
- Не поддерживается `JOIN` двух потоков (временное архитектурное ограничение).
- Не поддерживается обогащение из таблиц YDB — потоковые запросы могут использовать только [внешние таблицы S3](https://ydb.tech/docs/ru/concepts/query_execution/federated_query/s3/external_table.md?version=v26.1) для обогащения данными (поддержка обогащения из таблиц  YDB планируется к реализации в версии 26.1).
- Не поддерживается чтение и запись локальных топиков — для работы с ними необходимо использовать [внешние источники данных](https://ydb.tech/docs/ru/concepts/datamodel/external_data_source.md?version=v26.1), указывающие на текущую базу данных (ограничение планируется снять в версии 26.1).

{% endnote %}

Для работы с локальными топиками нужно использовать [внешние источники данных](https://ydb.tech/docs/ru/concepts/datamodel/external_data_source.md?version=v26.1):

- Создайте [внешний источник данных](https://ydb.tech/docs/ru/concepts/datamodel/external_data_source.md?version=v26.1), указывающий на ту же или другую базу данных YDB, и обращайтесь к топикам через них.

Также не поддерживаются в текущей версии:

- Признак [важного читателя](https://ydb.tech/docs/ru/concepts/datamodel/topic.md?version=v26.1#important-consumer) для потребителей, используемых потоковыми запросами.
- [Автопартиционирование](https://ydb.tech/docs/ru/concepts/datamodel/topic.md?version=v26.1#autopartitioning) (split/merge партиций) топиков, используемых потоковыми запросами. При увеличении числа партиций в топике, из которого читает работающий потоковый запрос, новые партиции не будут обрабатываться.
- Изменение текста работающего запроса. Чтобы изменить запрос, его нужно пересоздать — удалить с помощью [DROP STREAMING QUERY](https://ydb.tech/docs/ru/yql/reference/syntax/drop-streaming-query.md?version=v26.1) и создать заново с помощью [CREATE STREAMING QUERY](https://ydb.tech/docs/ru/yql/reference/syntax/create-streaming-query.md?version=v26.1).

## Управление {#control}

Потоковые запросы создаются, изменяются и удаляются с помощью YQL-команд:

- [CREATE STREAMING QUERY](https://ydb.tech/docs/ru/yql/reference/syntax/create-streaming-query.md?version=v26.1) — создание;
- [ALTER STREAMING QUERY](https://ydb.tech/docs/ru/yql/reference/syntax/alter-streaming-query.md?version=v26.1) — изменение и управление состоянием;
- [DROP STREAMING QUERY](https://ydb.tech/docs/ru/yql/reference/syntax/drop-streaming-query.md?version=v26.1) — удаление.

Состояние запросов доступно в системной таблице [`.sys/streaming_queries`](https://ydb.tech/docs/ru/dev/system-views.md?version=v26.1#streaming).

## Язык запросов {#syntax}

Потоковые запросы пишутся на [YQL](https://ydb.tech/docs/ru/yql/reference/index.md?version=v26.1) и поддерживают привычные SQL-конструкции: [SELECT](https://ydb.tech/docs/ru/yql/reference/syntax/select/index.md?version=v26.1), [WHERE](https://ydb.tech/docs/ru/yql/reference/syntax/select/where.md?version=v26.1), [GROUP BY](https://ydb.tech/docs/ru/yql/reference/syntax/select/group-by.md?version=v26.1), [JOIN](https://ydb.tech/docs/ru/yql/reference/syntax/select/join.md?version=v26.1). Для работы с временными окнами используется [GROUP BY HOP](https://ydb.tech/docs/ru/yql/reference/syntax/select/group-by.md?version=v26.1#group-by-hop).

Один потоковый запрос может читать несколько входных топиков, использовать конструкцию [UNION ALL](https://ydb.tech/docs/ru/yql/reference/syntax/select/union.md?version=v26.1#union-all) для объединения потоков данных и писать результат в несколько выходных топиков и/или таблиц (см. подробнее в статьях [Форматы при записи данных](https://ydb.tech/docs/ru/dev/streaming-query/streaming-query-formats.md?version=v26.1#write_formats) и [Запись в таблицы](https://ydb.tech/docs/ru/dev/streaming-query/table-writing.md?version=v26.1)).

## См. также

- [Гарантии доставки данных](https://ydb.tech/docs/ru/dev/streaming-query/guarantees.md?version=v26.1) — гарантии, аномалии при оконной агрегации и рекомендации;
- [Быстрый старт: чтение и запись в топики](https://ydb.tech/docs/ru/recipes/streaming_queries/topics.md?version=v26.1) — пошаговое руководство;
- [Форматы данных при чтении/записи из топиков](https://ydb.tech/docs/ru/dev/streaming-query/streaming-query-formats.md?version=v26.1) — форматы данных.
