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

<!-- source: ru/_includes/flow/concepts/pipeline-object.md -->
# Объект Pipeline в YTsaurus

В данном разделе описан Cypress-объект `pipeline` &mdash; единица деплоя YTsaurus Flow. Здесь рассматриваются модель данных, способы создания, внутренние таблицы и инструменты управления.

Связанные разделы: [Pipeline в глоссарии](https://ytsaurus.tech/docs/ru/flow/concepts/glossary.md#pipeline), [Spec, DynamicSpec и Config](https://ytsaurus.tech/docs/ru/flow/concepts/spec.md), [Stateful-обработка](https://ytsaurus.tech/docs/ru/flow/concepts/stateful.md), [Внутренние таблицы пайплайна](https://ytsaurus.tech/docs/ru/flow/concepts/glossary.md#inner-pipeline-tables).

## Модель данных { #data_model }

*Пайплайном (pipeline)* в YTsaurus называется специальный Cypress-объект типа `pipeline`. По устройству это map-узел &mdash; директория, в которой автоматически создаётся набор служебных [внутренних таблиц](#internal_tables), необходимых для работы Flow.

## Внутренние таблицы { #internal_tables }

<small>Таблица 1 &mdash; Внутренние таблицы пайплайна</small>

| Имя таблицы                         | Назначение                                                                                                                                                                                                                                                                                                        |
|-------------------------------------|-------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------|
| `input_messages`                    | Индекс входных сообщений [transform computation'ов](computation.md#ttransformcomputation), использующийся для дедупликации                                                                                                                                                                                        |
| `compact_input_messages`            | Компактный индекс входных сообщений. Используется по умолчанию для всех computation'ов, кроме случая, когда включён `experimental_enable_non_uint_key`; поведение переопределяется параметром `use_compact_input_messages` в [TComputationSpec](../generated_docs/all_yson_structs.md#NYT_NFlow_TComputationSpec) |
| `compact_output_messages`           | Не используется                                                                                                                                                                                                                                                                                                   |
| `compact_partition_output_messages` | Выходные сообщения transform computation'ов, физически сгруппированные по партициям и чанкованные по `stream_id` / `chunk_id` для оптимального чтения |
| `states`                            | Пользовательские и служебные [стейты](stateful.md), сохраняемые по [ключу](glossary.md#key)                                                                                                                                                                                                                       |
| `partition_states`                  | Стейты сохраняемые по партиции                                                                                                                                                                                                                                                                                    |
| `timers`                            | [Таймеры](glossary.md#timer) пользовательского кода                                                                                                                                                                                                                                                               |
| `controller_logs`                   | Логи событий [Controller](glossary.md#controller) в категории PublicFlowController                                                                                                                                                                                                                                |
| `flow_state`                        | Текущее состояние Flow                                                                                                                                                                                                                                                                                            |
| `flow_state_obsolete`               | KV-хранилище Flow для именнованных объектов (spec, dynamic_spec...)                                                                                                                                                                                                                                               |
| `flow_control` | Опубликованный адрес ведущего Controller; для Dyntable-бэкенда выборов также аренда лидерства |
| `key_visitor_states` | Сохраняемый курсор обхода диапазонов ключей и состояние прохода |
| `partition_transactions`            | Служебная таблица для безопасного ретрая транзакций                                                                                                                                                                                                                                                               |
| `leases` | Для Dyntable-бэкенда выборов: владельцы аренд партиций и общий для пайплайна срок действия аренд, защищающий от записи устаревшими воркерами |
| `leader_election_lock` | Блокировка выборов Controller при использовании Chaos-бэкенда |

После создания пайплайна они появляются под путём `<pipeline_path>/<table_name>` и автоматически монтируются.

Состав таблиц, их схемы и базовые атрибуты для `create pipeline` заданы в [`definitions.yson`](https://github.com/ytsaurus/ytsaurus/tree/main/yt/yt/flow/library/pipeline_tables/definitions.yson).


Библиотека `yt_sync_mini` дополнительно применяет [пресеты физических атрибутов](https://github.com/ytsaurus/ytsaurus/tree/main/yt/yt/flow/library/python/pipeline_tables/presets.py).


Таблицы `input_messages` и `compact_input_messages` не накапливают историю: записи о входных сообщениях нужны только для дедупликации (exactly-once) и удаляются после обработки. Controller продвигает [SystemWatermark](https://ytsaurus.tech/docs/ru/flow/concepts/watermarks.md) этих таблиц, и строки с уже обработанными сообщениями физически удаляются механизмом очистки динамических таблиц &mdash; в том числе после завершения пайплайна. Поэтому пустая таблица `input_messages` &mdash; это штатная очистка, а не потеря данных.

{% note warning "Внимание" %}

Внутренние таблицы являются служебными и их структура может меняться между релизами Flow. Не следует читать или писать в них напрямую из пользовательского кода. Для чтения отладочных данных (например, логов контроллера) используйте `yt flow show-logs` и другие команды семейства `yt flow`.

{% endnote %}

## Создание пайплайна { #create }


### Через yt_sync_mini { #yt-sync-mini }

Рекомендуемый способ создания пайплайна в опенсорсе &mdash; Python-библиотека [`yt_sync_mini`](https://github.com/ytsaurus/ytsaurus/tree/main/yt/yt/flow/library/python/yt_sync_mini). Она создаёт map-узел типа `pipeline`, все [внутренние таблицы](#internal_tables) с корректными схемами и физическими атрибутами и сразу монтирует их. Операция идемпотентна &mdash; повторный запуск над уже существующим пайплайном является no-op.

```python
import yt.wrapper as yt

from yt.yt.flow.library.python.yt_sync_mini import create_pipeline

client = yt.YtClient(proxy="<cluster>")
create_pipeline(client, "<pipeline_path>")
```

### Низкоуровневое создание Cypress-узла { #low-level-create }

Для интеграции в собственную систему деплоя пайплайн можно создать штатным механизмом `create` &mdash; так же, как и другие типы Cypress-объектов (table, map_node, queue_consumer и т. д.). По умолчанию `create pipeline` создаёт и монтирует внутренние таблицы. Атрибут `initialize_tables=%false` передают при создании объекта, чтобы отключить эту инициализацию; тогда пользователь сам отвечает за создание и монтирование таблиц с корректными схемами и атрибутами.

#### Через YTsaurus CLI

```bash
yt --proxy <cluster> create pipeline <pipeline_path>
```

Если внутренние таблицы будут созданы отдельно:

```bash
yt --proxy <cluster> create pipeline <pipeline_path> --attributes '{initialize_tables=%false}'
```

#### Через Python (YTsaurus wrapper)

```python
import yt.wrapper as yt

client = yt.YtClient(proxy="<cluster>")
client.create(
    "pipeline",
    "<pipeline_path>"
)
```

#### Через C++ (YTsaurus native client)

```cpp
#include <yt/yt/flow/library/cpp/native_client/pipeline_init.h>

NYT::NApi::TCreateNodeOptions options;

auto nodeId = NYT::NFlow::CreatePipelineNode(client, pipelinePath, options);
```


## External State { #external-state }

При смене схемы [внутренних таблиц](#internal_tables) в новой версии Flow апгрейд формата выполняется отдельной миграцией &mdash; см. [Внутренние таблицы пайплайна](https://ytsaurus.tech/docs/ru/flow/concepts/glossary.md#inner-pipeline-tables) и [Базовые правила выкатки](https://ytsaurus.tech/docs/ru/flow/devops/vanilla/releases.md#release-and-configure-basic-rules).

Если пайплайн использует [External State](https://ytsaurus.tech/docs/ru/flow/concepts/stateful.md) (пользовательские таблицы за пределами узла), их создание и эволюция схем &mdash; ответственность пользователя. Операции выполняются стандартными командами `yt create table ... --attributes '{dynamic=%true; schema=...}'` и `yt mount-table` &mdash; см. примеры в разделе [Команда create](https://ytsaurus.tech/docs/ru/user-guide/storage/cypress-example.md#create).

## См. также { #see_also }

- [Глоссарий: Pipeline](https://ytsaurus.tech/docs/ru/flow/concepts/glossary.md#pipeline)
- [Внутренние таблицы пайплайна](https://ytsaurus.tech/docs/ru/flow/concepts/glossary.md#inner-pipeline-tables)
- [Базовые правила выкатки пайплайна](https://ytsaurus.tech/docs/ru/flow/devops/vanilla/releases.md#release-and-configure-basic-rules)
<!-- endsource: ru/_includes/flow/concepts/pipeline-object.md -->
