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

<!-- source: ru/_includes/flow/concepts/glossary.md -->
# Глоссарий YTsaurus Flow

## Введение {#introduction}

В данной статье собраны все ключевые понятия YTsaurus Flow: от базовой модели обработки данных до архитектуры и механизмов обеспечения корректности. Материал выстроен от простого к сложному: сначала рассматривается потоковая обработка, а затем &mdash; параллелизм, внешние подключения, устройство системы и управление пайплайном.

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

## Основная модель {#basic-model}

Для понимания Flow как системы потоковой обработки данных можно провести аналогию с заводским конвейером:

- [Пайплайн](#pipeline) — целый завод, решающий конкретную бизнес-задачу.
- [Стрим](#stream) — лента конвейера, по которой движутся данные.
- [Сообщение](#message) — отдельный элемент на ленте.
- [Компьютейшен](#computation) — рабочая станция, которая берёт элементы с одной или нескольких лент, обрабатывает их и кладёт результаты на другие ленты.

### Pipeline (пайплайн) {#pipeline}

Это конкретная выполняемая бизнес-задача. Может содержать множество различных узлов группировки и обработки данных.

С точки зрения YTsaurus пайплайн представлен Cypress-объектом типа `pipeline`. Подробности про устройство объекта, способы создания и управления &mdash; в разделе [Объект Pipeline](https://ytsaurus.tech/docs/ru/flow/concepts/pipeline-object.md).


### Stream (стрим) {#stream} {#stream-and-computation}

Стрим — схематизированный именованный поток данных, соединяющий [компьютейшены](#computation) в графе пайплайна. Каждый стрим имеет фиксированную схему (набор типизированных полей) и идентифицируется по `stream_id`. Стрим состоит из [сообщений](#message). У каждого стрима ровно один писатель — [компьютейшен](#computation), который выводит в него данные — и произвольное количество читателей.

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

### Computation (компьютейшен) {#computation}

Компьютейшен — узел графа пайплайна, выполняющий конкретное преобразование данных. Читает [сообщения](#message) из входных [стримов](#stream), обрабатывает их и записывает результаты в выходные [стримы](#stream). Может иметь как множество входов, так и множество выходов. Входной поток разбивается на [партиции](#partition) для параллельной обработки.

Подробнее в разделе [Computation](https://ytsaurus.tech/docs/ru/flow/concepts/computation.md).

### Passthrough Computation (passthrough-компьютейшен) {#passthrough}

Вид [компьютейшена](#computation), не содержащий пользовательской бизнес-логики: входящие сообщения конвертируются в схему выходного [стрима](#stream) и передаются дальше без изменений. Используется для простого приведения схем между стримами.

Подробнее — [Computation](https://ytsaurus.tech/docs/ru/flow/concepts/computation.md#passthrough).

### Swift {#swift}

Принцип обработки данных во Flow, при котором результат работы [компьютейшена](#computation) не сохраняется в YTsaurus. Вместо этого функция преобразования должна быть строго детерминированной — при необходимости результат вычисляется повторно. Это снижает нагрузку на YTsaurus при сохранении гарантий [exactly-once](#exactly-once).

Подробнее в разделе [Swift](https://ytsaurus.tech/docs/ru/flow/concepts/swift.md).

### Message (сообщение) {#message}

Одно сообщение в рамках потока.

В контексте одного [компьютейшена](#computation) бывает несколько видов сообщений:
- `input` &mdash; сообщения от других [компьютейшенов](#computation), по сути, `output` другого [компьютейшена](#computation), но со сгруппированной схемой.
- `source` &mdash; сообщения из внутренних [сорсов](#source).
- `timer` &mdash; специальные внутренние сообщения-[таймеры](#timer).
- `output` &mdash; результат работы [компьютейшена](#computation), который становится публичным и может попадать во внутренний [синк](#sink) и в другие [компьютейшены](#computation).

### Key (ключ группировки) {#key}

Набор значений полей [сообщения](#message), определяемых через `group_by_schema` в [спеке](https://ytsaurus.tech/docs/ru/flow/concepts/spec.md) [компьютейшена](#computation). Все сообщения с одинаковым ключом направляются в одну [партицию](#partition) — каждой партиции назначен диапазон ключей `[LowerKey; UpperKey)`. Ключ является единицей изоляции для [стейта](#state), [таймеров](#timer) и [гарантий порядка](https://ytsaurus.tech/docs/ru/flow/concepts/ordering.md#ordering-guarantees).

Понятие ключа зависит от типа сообщения:
- Для сообщений типа `input` — ключ определяется через `group_by_schema`.
- Для сообщений типа `source` — ключ является партиционным ключом источника.
- Для сообщений типа `output` — ключа нет.

### Lineage (родословная) {#lineage}

Совокупность входных [сообщений](#message) и [таймеров](#timer), из которых были получены конкретные выходные результаты [компьютейшена](#computation). Фреймворк использует lineage для вычисления системных выходных сообщений и для обеспечения [гарантий порядка](https://ytsaurus.tech/docs/ru/flow/concepts/ordering.md#ordering-guarantees) производных сообщений.

Подробнее в разделе [Lineage](https://ytsaurus.tech/docs/ru/flow/concepts/lineage.md).

## Timer (таймер) {#timer}

Механизм отложенного вызова, привязанный к конкретному [ключу группировки](#key) в рамках `TTransformComputation`. Таймер позволяет [компьютейшену](#computation) сказать: «разбуди меня, когда время X наступит».

Каждый таймер содержит два временных поля:
- `TriggerTimestamp` — момент срабатывания. Когда [EventWatermark](https://ytsaurus.tech/docs/ru/flow/concepts/glossary.md#timestamps-and-watermarks) (или другой сконфигурированный вотермарк) превысит это значение, таймер передаётся в компьютейшен для обработки.
- `EventTimestamp` — бизнес-время исходного события, «придержанное» в таймере.

Таймеры надёжно сохраняются в YTsaurus. Результат обработки таймера фиксируется в той же транзакции, в которой таймер удаляется, — это обеспечивает гарантии [exactly-once](#exactly-once).

{% note warning %}

В текущей реализации все активные таймеры дополнительно хранятся в памяти процесса. При большом количестве таймеров это может приводить к ошибкам `Out of Memory` при старте джобов.

{% endnote %}

Подробнее в разделе [Таймеры](https://ytsaurus.tech/docs/ru/flow/concepts/timers.md).

## Resource (ресурс) {#resource}

Ресурсы существуют для описания данных, общих для нескольких [джобов](#job). Это могут быть клиенты к YTsaurus и другим системам, кеши для обращений во внешние системы, модели машинного обучения и т. п.

На данный момент ресурсы могут быть только статическими (то есть неизменяемыми без перезапуска).

## Distributed Throttler (распределённый троттлер) {#distributed-throttler}

Во Flow есть возможность создавать именованные распределённые троттлеры — общие для всех [джобов](#job) пайплайна. Их используют для ограничения нагрузки на внешние API, выравнивания пропускной способности между [партициями](#partition) одного [компьютейшена](#computation) или притормаживания чтения из источников. Реализованы как [token bucket](https://en.wikipedia.org/wiki/Token_bucket) на контроллере; джобы запрашивают квоту перед обработкой сообщений или вручную из пользовательского кода. Подробнее в разделе [Distributed Throttler](https://ytsaurus.tech/docs/ru/flow/concepts/distributed_throttler.md).

## State (стейт) {#state}

Персистентные данные, связанные с конкретным [ключом группировки](#key) в рамках [компьютейшена](#computation). Хранятся в динамических таблицах YTsaurus и обновляются атомарно при коммите [эпохи](#epoch). Пустое значение стейта соответствует отсутствию строки в таблице.

Подробнее в разделе [Stateful processing](https://ytsaurus.tech/docs/ru/flow/concepts/stateful.md).

## Параллелизм и исполнение {#parallelism-and-execution}

Для обеспечения масштабируемости процесса обработки данных, Flow разбивает каждый [компьютейшен](#computation) на множество [партиций](#partition), которые обрабатываются параллельно в рамках [джобов](#job).

### Partition (партиция) {#partition}

Входной поток в конкретный [компьютейшен](#computation) может быть достаточно велик, поэтому для целей параллельной обработки данных [компьютейшен](#computation) разбивается на множество партиций. Ориентировочные ограничения по потоку для одной партиции: не более 1 MB/s и не более 1000 сообщений в секунду.

За счёт поля `group_by_schema` в [спеке](https://ytsaurus.tech/docs/ru/flow/concepts/spec.md) [компьютейшена](#computation) входной поток группируется и данные по одному [ключу](#key) всегда будут попадать в конкретную партицию. Для этого за каждой партицией определён диапазон ключей `[LowerKey; UpperKey)`.

### Job (джоб) {#job}

Конкретный запуск обработчика конкретной [партиции](#partition) на [воркере](#worker), выбранном [контроллером](#controller).

### Epoch (эпоха) {#epoch}

Джоб обрабатывает данные эпохами для регулярного надёжного сохранения прогресса обработки данных. Механизм эпох позволяют поддерживать гарантии [exactly-once](#exactly-once) при разумной нагрузке на YTsaurus, а также для минимизации времени простоя при восстановлении от сбоев.

Некоторые эпохи могут не содержать коммитов в YTsaurus. Некоторые реализации [компьютейшена](#computation) могут вообще не делать коммиты.

### Layout {#layout}

Список всех [джобов](#job) и [партиций](#partition) в системе. Иными словами &mdash; описание, что где сейчас работает.

## Внешние подключения {#external-connections}

Flow взаимодействует с внешними системами (очередями, таблицами и т. д.) через [коннекторы](#connector). Каждый коннектор предоставляет [сорс](#source) для чтения и/или [синк](#sink) для записи.

### Connector (коннектор) {#connector}

Компонент, обеспечивающий связь пайплайна с некоторой внешней системой через чтение сообщений из [source](#source) коннектора и/или запись сообщений в [sink](#sink) коннектора. Коннектор может включать в себя и [source](#source), и [sink](#sink) (например, в случае [QYT](https://ytsaurus.tech/docs/ru/flow/connectors/queue.md)) или же только [source](#source) (например, в случае [servicelog](https://ytsaurus.tech/docs/ru/flow/connectors/servicelog.md)).

В коде и в спеках описываются напрямую source'ы и sink'и. Сам коннектор не существует как сущность, это скорее элемент логической группировки source'ов и sink'ов по внешней системе, с которой идёт взаимодействие.

Подробнее о коннекторах можно прочитать в [документации по коннекторам](https://ytsaurus.tech/docs/ru/flow/connectors/about.md).

### Source (сорс) {#source}

Компонент, позволяющий [компьютейшену](#computation) читать сообщения из внешней системы. Реализуется в виде отдельного класса и конфигурируется в спеке [компьютейшена](#computation).

### Sink (синк) {#sink}

Компонент, позволяющий [компьютейшену](#computation) записывать сообщения во внешнюю систему. Реализуется в виде отдельного класса и конфигурируется в спеке [компьютейшена](#computation).

## Архитектура системы {#system-architecture}

Запущенный пайплайн состоит из набора процессов двух ролей — [контроллера](#controller) и [воркера](#worker). Контроллер управляет распределением работы, а воркеры непосредственно обрабатывают данные.

### Controller (контроллер) {#controller}

Управляет всем пайплайном. Назначает [джобы](#job) [воркерам](#worker), предоставляет API для управления пайплайном, отвечает за синхронизацию всех настроек, следит за состоянием всех [воркеров](#worker) и собирает от них статусы. Может быть запущен в нескольких экземплярах (рекомендуемое &mdash; 2-3), тогда один из инстансов будет лидером, а остальные будут ожидать в режиме stand-by для обеспечения отказоустойчивости.

Спеки для выполнения задаются в контроллер через специальный API.

#### Computation Controller {#computation-controller}

Контроллер конкретного [компьютейшена](#computation): отвечает за создание необходимого количества [партиций](#partition) с нужными настройками. Не путать с просто [контроллером](#controller). Computation Controller — это класс, а контроллер — это программа, которая управляет работой кластера и создаёт необходимые Computation Controller-ы.

#### JobManager {#job-manager}

Модуль, отвечающий за создание и распределение [джобов](#job) по [воркерам](#worker).

#### LeaseManager {#lease-manager}

Модуль, отвечающий за создание и пинг `Lease` &mdash; мастер-транзакций, которые используются джобами при коммитах в качестве `prerequisite_transaction_ids`. Позволяют обеспечить гарантии [exactly-once](#exactly-once).

### Worker (воркер) {#worker}

Один из двух основных видов процессов в системе, осуществляющий сами вычисления в рамках назначенных ему [джобов](#job). Регулярно шлёт [хартбиты](#heartbeat) в [контроллер](#controller), отправляя ему статусы всех джобов. В ответ получает актуальные настройки системы, в частности [Layout](#layout).

#### Хартбит (heartbeat) {#heartbeat}

Периодическое сообщение, которое [воркер](#worker) отправляет [контроллеру](#controller) со статусами всех своих [джобов](#job). В ответе контроллера приходят актуальные настройки системы, в частности [Layout](#layout); воркер, переставший слать хартбиты, считается недоступным.

#### Message Distributor {#message-distributor}

Модуль [воркера](#worker), отвечающий за рассылку сообщений между различными [джобами](#job). Выполняется на базе текущего [лэйаута](#layout) системы. Шлёт сообщение до тех пор, пока не получит от получателя подтверждения обработки, включая факт надёжного коммита результата обработки (в кодовой базе упоминается как `MarkPersisted`).

#### Resource Manager {#resource-manager}

Компонент [воркера](#worker), создающий [ресурсы](#resource) для общего пользования. Например, несколько [компьютейшенов](#computation) могут использовать одну и ту же бинарную базу или один и тот же клиент YTsaurus.

#### BufferStateManager {#buffer-state-manager}

Компонент, отвечающий за создание и поддержание актуальных размеров всех (или почти всех) буферов в системе. У каждой джобы есть буферы входящих и исходящих сообщений. Задача `BufferStateManager` &mdash; обеспечить необходимый для бесперебойной работы размер буферов.

### Companion (компаньон) {#companion}

Отдельный процесс, запускаемый на одном хосте с [воркером](#worker) и выполняющий пользовательский код на [Python](https://ytsaurus.tech/docs/ru/flow/python/getting-started.md), [Go](https://ytsaurus.tech/docs/ru/flow/go/getting-started.md), [Java или Kotlin](https://ytsaurus.tech/docs/ru/flow/java/getting-started.md) либо C++. Взаимодействует с [воркером](#worker) по gRPC. Выполняет [компьютейшены](#computation) вне процесса воркера: как на языках, отличных от C++, так и с помощью [C++-компаньона](https://ytsaurus.tech/docs/ru/flow/concepts/companion.md#cpp-companion).

Подробнее в разделе [Companion](https://ytsaurus.tech/docs/ru/flow/concepts/companion.md).

## Время и корректность {#timestamps-and-watermarks}

Корректная обработка событий во времени — одна из ключевых задач потоковой системы. Flow отслеживает три вида временных меток у каждого сообщения и поддерживает механизм вотермарков для определения прогресса обработки.

### EventTimestamp, SystemTimestamp, AlignmentTimestamp, StabilizedEventTimestamp и Watermarks

В отношении каждого сообщения во Flow могут использоваться следующие таймстемпы:
* [EventTimestamp](https://ytsaurus.tech/docs/ru/flow/concepts/watermarks.md#eventtimestamp) &mdash; истинное время события, ассоциированного с данным сообщением.
* [SystemTimestamp](https://ytsaurus.tech/docs/ru/flow/concepts/watermarks.md#systemtimestamp) &mdash; время создания конкретного сообщения. Для сообщений, сгенерированных внутри Flow, в качестве `SystemTimestamp` берётся «время YTsaurus». Для `source` потоков &mdash; время появления сообщения в удалённой системе.
* [AlignmentTimestamp](https://ytsaurus.tech/docs/ru/flow/concepts/ordering.md#alignment-timestamp) &mdash; таймстемп, используемый для выравнивания прогресса обработки партиций. Вычисляется автоматически.
* [StabilizedEventTimestamp](https://ytsaurus.tech/docs/ru/flow/concepts/ordering.md#stabilized-event-timestamp) &mdash; таймстемп, вычисляющийся на основе `AlignmentTimestamp`, используется как неубывающее на потоке сообщений приближение `EventTimestamp` (при совпадающем наборе [ключей](#key) в [lineage](#lineage)). 

Первые три таймстемпа хранятся в самом сообщении, а `StabilizedEventTimestamp` вычисляется на лету при необходимости.

[Вотермарк (watermark)](https://ytsaurus.tech/docs/ru/flow/concepts/watermarks.md) &mdash; это временная метка, означающая, что система больше не ожидает поступления более старых событий. На каждом стриме в системе поддерживаются текущие `SystemWatermark` и `EventWatermark`.

{% note info %}

К сожалению, обычные протоколы записи данных в очереди не могут гарантировать, что в очередь не будет записано вчерашнее событие. Поэтому про `SystemWatermark` по внутренним потокам можно говорить, что он является абсолютным, а про `EventWatermark` &mdash; лишь то, что он эвристичен и является приближённым.

{% endnote %}

## Обеспечение Exactly Once {#exactly-once}

Flow по умолчанию обеспечивает exactly-once семантику обработки событий. Механизм основан на барьерных Lease-транзакциях, дедупликации входных сообщений по `message_id` (для внутренних потоков) и по оффсетам (для источников), а также на атомарности выходных данных через транзакции [эпохи](https://ytsaurus.tech/docs/ru/flow/concepts/glossary.md#epoch). При необходимости семантику можно ослабить до at-least-once или at-most-once.

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

## Дедупликация сообщений и порядок обработки сообщений {#deduplication}

Основным способом дедупликации сообщений является дедупликация по `message_id`, а не по оффсету. Механизм с оффсетами используется только при работе с очередями. В том числе, Flow не даёт никакой гарантии на порядок обработки событий.

## Управление пайплайном {#manage-pipeline}

Этот раздел описывает операции жизненного цикла пайплайна: от описания его конфигурации через спеки до релиза, миграции и работы в различных окружениях.

### Spec и DynamicSpec {#spec-and-dynamic-spec}

Любой выполняемый на Flow пайплайн представляет собой комбинацию из двух YSON конфигов, называемых `Spec` и `DynamicSpec`.

- `Spec` &mdash; определяет статические свойства пайплайна, его топологию, свойства узлов, связи между ними, типы объектов и т. п. Статическая спека может меняться только при условии остановки пайплайна.
- `DynamicSpec` &mdash; определяет динамические свойства пайплайна.

Соответственно, спеки содержат свойства всех компонент системы, как пользовательского слоя (например, список [компьютейшенов](#computation) и схемы [стримов](#stream) между ними), так и настроек системного уровня (кодеки сжатия в системных таблицах, размеры буферов и т. п.).

### Релиз пайплайна {#release-pipeline}

Обновление пайплайна, которое включает изменение его исполняемых файлов.

### Обновление спек пайплайна {#update-pipeline-specs}

Обновление [статической](*static_spec_upd) и/или динамической [спеки](https://ytsaurus.tech/docs/ru/flow/concepts/spec.md) пайплайна.

### Запуск и остановка пайплайна {#start-stop-pause-pipeline}

Изменение целевого состояния пайплайна и переход пайплайна в это состояние. Существующие состояния:

- `unknown` — пайплайн ещё не запущен.
- `working` — пайплайн работает, сообщения обрабатываются.
- `stopped` — пайплайн остановлен, все сообщения дообработаны.
- `paused` — пайплайн приостановлен (джобы остановлены, промежуточные сообщения могут быть не дообработаны).
- `draining` — переходное состояние: пайплайн в процессе остановки.
- `pausing` — переходное состояние: пайплайн в процессе приостановки.
- `completed` — финальное состояние: все источники пайплайна были [конечными](https://ytsaurus.tech/docs/ru/flow/python/testing.md) (`finite = true`) и все сообщения из них обработаны. Выйти из этого состояния нельзя — потребуется пересоздать пайплайн. Чаще всего встречается в [интеграционных тестах](https://ytsaurus.tech/docs/ru/flow/python/testing.md), а также может возникать в продакшн-пайплайнах при некорректной последовательности действий во время деплоя (например, если источник был ошибочно помечен как конечный) или при выставлении некорректной спеки.

### Внутренние таблицы пайплайна {#inner-pipeline-tables}

Служебные таблицы, которые размещаются в директории пайплайна в YTsaurus. Они необходимы для работы пайплайна (например, для дедупликации входных сообщений, хранения информации по партициям и т. д.).

Полный список внутренних таблиц и их назначение приведены в разделе [Объект Pipeline &rarr; Внутренние таблицы](https://ytsaurus.tech/docs/ru/flow/concepts/pipeline-object.md#internal_tables).

### Пользовательские таблицы {#user-tables}

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

- Упорядоченные таблицы &mdash; очереди, из которых пользовательский пайплайн читает входные события или в которые пишет выходные.
- Сортированные таблицы с состояниями сущностей (профили баннера/оффера/пользователя и т. д.), которые пользователь строит в процессе работы пайплайна.

### Миграция {#migration}

Процесс изменения структуры объектов в YTsaurus &mdash; например, добавление новых таблиц, изменение схемы существующих таблиц, их удаление.

Основные причины для выполнения миграций:

* Выкатка новой версии пайплайна, не совместимой с текущим состоянием объектов в YTsaurus. Миграция позволяет без потери данных обновить объекты до состояния, с которым сможет работать новая версия пайплайна.
* Изменение типа таблицы (replicated/standalone/chaos) из соображений производительности или доступности.

Примеры изменений:

#|
|| **Требуют миграции**

- Изменение множества таблиц (добавление новых или удаление старых).
- Изменение схем таблиц.
- Изменение типа таблиц (replicated/standalone/chaos).

| **Не требуют миграции**

- Изменение большинства атрибутов таблиц (например, изменение настроек компактификации, влияющее только на производительность).
- Изменение списка реплик.

||
|#


### Окружение (environment) {#environment}

Инсталляция, где запускается пайплайн, &mdash; тестинг, препродакшн или продакшн.

Является синонимом понятия «стейдж» (stage).





## См. также

- [Computation](https://ytsaurus.tech/docs/ru/flow/concepts/computation.md)
- [Spec и DynamicSpec](https://ytsaurus.tech/docs/ru/flow/concepts/spec.md)
- [Watermarks и Timers](https://ytsaurus.tech/docs/ru/flow/concepts/watermarks.md)
- [Stateful processing](https://ytsaurus.tech/docs/ru/flow/concepts/stateful.md)
- [Коннекторы](https://ytsaurus.tech/docs/ru/flow/connectors/about.md)
<!-- endsource: ru/_includes/flow/concepts/glossary.md -->


[*static_spec_upd]: Можно менять только при условии остановки пайплайна.