Объект Pipeline в YTsaurus
В данном разделе описан Cypress-объект pipeline — единица деплоя YTsaurus Flow. Здесь рассматриваются модель данных, способы создания, внутренние таблицы и инструменты управления.
Связанные разделы: Pipeline в глоссарии, Spec, DynamicSpec и Config, Stateful-обработка, Внутренние таблицы пайплайна.
Модель данных
Пайплайном (pipeline) в YTsaurus называется специальный Cypress-объект типа pipeline. По устройству это map-узел — директория, в которой автоматически создаётся набор служебных внутренних таблиц, необходимых для работы Flow.
Внутренние таблицы
Таблица 1 — Внутренние таблицы пайплайна
| Имя таблицы | Назначение |
|---|---|
input_messages |
Индекс входных сообщений transform computation'ов, использующийся для дедупликации |
compact_input_messages |
Компактный индекс входных сообщений. Используется по умолчанию для всех computation'ов, кроме случая, когда включён experimental_enable_non_uint_key; поведение переопределяется параметром use_compact_input_messages в TComputationSpec |
compact_output_messages |
Не используется |
compact_partition_output_messages |
Выходные сообщения transform computation'ов, физически сгруппированные по партициям и чанкованные по stream_id / chunk_id для оптимального чтения |
states |
Пользовательские и служебные стейты, сохраняемые по ключу |
partition_states |
Стейты сохраняемые по партиции |
timers |
Таймеры пользовательского кода |
controller_logs |
Логи событий 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.
Библиотека yt_sync_mini дополнительно применяет пресеты физических атрибутов.
Таблицы input_messages и compact_input_messages не накапливают историю: записи о входных сообщениях нужны только для дедупликации (exactly-once) и удаляются после обработки. Controller продвигает SystemWatermark этих таблиц, и строки с уже обработанными сообщениями физически удаляются механизмом очистки динамических таблиц — в том числе после завершения пайплайна. Поэтому пустая таблица input_messages — это штатная очистка, а не потеря данных.
Внимание
Внутренние таблицы являются служебными и их структура может меняться между релизами Flow. Не следует читать или писать в них напрямую из пользовательского кода. Для чтения отладочных данных (например, логов контроллера) используйте yt flow show-logs и другие команды семейства yt flow.
Создание пайплайна
Через yt_sync_mini
Рекомендуемый способ создания пайплайна в опенсорсе — Python-библиотека yt_sync_mini. Она создаёт map-узел типа pipeline, все внутренние таблицы с корректными схемами и физическими атрибутами и сразу монтирует их. Операция идемпотентна — повторный запуск над уже существующим пайплайном является no-op.
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-узла
Для интеграции в собственную систему деплоя пайплайн можно создать штатным механизмом create — так же, как и другие типы Cypress-объектов (table, map_node, queue_consumer и т. д.). По умолчанию create pipeline создаёт и монтирует внутренние таблицы. Атрибут initialize_tables=%false передают при создании объекта, чтобы отключить эту инициализацию; тогда пользователь сам отвечает за создание и монтирование таблиц с корректными схемами и атрибутами.
Через YTsaurus CLI
yt --proxy <cluster> create pipeline <pipeline_path>
Если внутренние таблицы будут созданы отдельно:
yt --proxy <cluster> create pipeline <pipeline_path> --attributes '{initialize_tables=%false}'
Через Python (YTsaurus wrapper)
import yt.wrapper as yt
client = yt.YtClient(proxy="<cluster>")
client.create(
"pipeline",
"<pipeline_path>"
)
Через C++ (YTsaurus native client)
#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
При смене схемы внутренних таблиц в новой версии Flow апгрейд формата выполняется отдельной миграцией — см. Внутренние таблицы пайплайна и Базовые правила выкатки.
Если пайплайн использует External State (пользовательские таблицы за пределами узла), их создание и эволюция схем — ответственность пользователя. Операции выполняются стандартными командами yt create table ... --attributes '{dynamic=%true; schema=...}' и yt mount-table — см. примеры в разделе Команда create.