Объект 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.

См. также

Предыдущая