Ресурсы
Ресурс — это объект, общий для всех джобов одного процесса Flow. Через ресурс вычисление может
получить клиент внешней системы, справочник, модель или другой объект, который дорого создавать
для каждого сообщения.
Ресурсы объявляются в секции resources спецификации пайплайна, а вычисления перечисляют нужные
им ресурсы в required_resource_ids. Один экземпляр ресурса создаётся на процесс воркера.
Конструктор должен быть лёгким; длительная инициализация выполняется в Load().
Примеры ниже показывают только относящиеся к ресурсам фрагменты классов и спецификации.
Полный проект пайплайна и команды сборки описаны в быстром старте.
Пользовательский ресурс обычно наследуется от TResourceBase и регистрируется макросом
YT_FLOW_DEFINE_RESOURCE:
class TMyResource
: public TResourceBase
{
public:
using TResourceBase::TResourceBase;
TFuture<void> Load(
const THashMap<TResourceId, IResourcePtr>& dependencies) override;
};
YT_FLOW_DEFINE_RESOURCE(TMyResource);
Параметры класса задаются в resources.<name>.parameters и доступны через GetParameters().
Зависимости от других ресурсов задаются в resources.<name>.dependencies.
Где загружать ресурс
Для каждого требования можно независимо указать, нужен ли ресурс на воркере и на контроллере:
required_resource_ids = {
geobase = {
worker = %true;
controller = %false;
};
};
Большинство ресурсов с пользовательскими данными нужны только на воркерах. Ресурс загружается на
контроллере, только если он нужен логике вычисления, которая действительно исполняется там.
Ограничение распространяется и на зависимости ресурса.
Контроллер ресурса
Контроллерная часть нужна ресурсу, если Flow должен следить за внешней версией данных и сообщать
воркерам, какую версию использовать. Для файлов эта логика уже реализована в
TResourceControllerBase.
Для дополнительной логики можно определить собственный контроллер:
class TMyResourceController
: public TResourceControllerBase
{
public:
using TResourceControllerBase::TResourceControllerBase;
protected:
NYTree::INodePtr DoBuildTargetRevisionSpec() override;
void DoCollectStatuses(
const THashMap<std::string, TWorkerResourceStatusPtr>& workerStatuses,
const TWorkerResourceStatusPtr& controllerStatus) override;
NYTree::IMapNodePtr DoGetView() override;
};
class TMyResource
: public TResourceBase
{
public:
using TController = TMyResourceController;
using TResourceBase::TResourceBase;
};
Контроллер публикует целевую ревизию ресурса. Воркеры получают её через обычную переконфигурацию
и сообщают, какую ревизию уже применили. Поэтому контроллер различает ситуацию «цель ещё не
доставлена» и ситуацию «цель доставлена, но ресурс ещё готовится».
DoCollectStatuses() получает состояния ресурса на живых воркерах, а DoGetView() формирует
часть представления ресурса в Flow view. Если собственное состояние должно переживать перезапуск
контроллера, его можно сохранить через контекст DoInit().
Ресурсы на основе файлов
Файловый ресурс подходит для справочников, моделей и других неизменяемых структур, которые можно
построить из одного или нескольких файлов. Провайдеры объявляются именованной картой
file_providers рядом с обычными параметрами ресурса:
resources = {
geobase = {
resource_class_name = "TGeobaseResource";
parameters = {
format = "binary";
};
file_providers = {
countries = {
file_provider_class_name = "NYT::NFlow::TYTFileProvider";
parameters = {
path = "<cluster=primary>//models/countries";
};
};
cities = {
file_provider_class_name = "NYT::NFlow::TYTFileProvider";
parameters = {
path = "<cluster=primary>//models/cities";
};
};
};
};
};
Такой ресурс должен загружаться только на воркерах: во всех достижимых требованиях укажите
controller = %false. Некорректная спецификация отклоняется до запуска пайплайна.
Реализация ресурса
Для стандартного сценария наследуйтесь от TFileResourceBase<TData> и реализуйте Initialize().
Если структура требует дополнительной проверки, переопределите Validate():
class TGeobaseResource
: public TFileResourceBase<TGeobase>
{
protected:
TGeobasePtr Initialize(
const TMaterializedFileProviderSnapshotPtr& files) override
{
const auto& countries = files->GetFileProvider(TFileProviderId("countries"));
const auto& cities = files->GetFileProvider(TFileProviderId("cities"));
return LoadGeobase(countries->GetRootPath(), cities->GetRootPath());
}
void Validate(const TGeobasePtr& geobase) override;
};
Flow обнаруживает поколение загрузки каждого именованного провайдера и собирает из них один снимок.
На каждом воркере снимок последовательно скачивается, инициализируется и проверяется. Новые данные
становятся доступны только после успеха всех трёх шагов для всех файлов. При ошибке ресурс
повторяет попытку через file_provider_update_retry_period. При первом запуске Load() остаётся
незавершённым до успешной попытки или получения новой цели. При обновлении ресурс продолжает
отдавать предыдущую исправную версию.
Контроллер хранит активный снимок и следующий подготавливаемый снимок. Новый воркер сначала
поднимает активный снимок, а затем занимается следующим. Контроллер допускает подготавливаемый
снимок после его успешной проверки хотя бы на одном актуальном экземпляре ресурса с текущей
целевой ревизией. Это проверка пригодности, а не одновременное переключение всех воркеров и не
кворум. После допуска воркеры переключаются постепенно и могут завершить переход в разное время,
но каждый из них получает один и тот же набор источников и поколений загрузки.
Для TYTFileProvider и TYTDirectoryLastFileProvider ключ кеша состоит из кластера, выбранного
пути, идентификатора объекта и content_revision. Изменение содержимого или замена объекта создаёт
новое поколение загрузки, которое проходит обычную подготовку и проверку пригодности. Изменение
остальных атрибутов не вызывает reload. Jobs и процессы воркеров продолжают работать; старые
accessor удерживают предыдущие данные до освобождения. file_snapshot_min_creation_period
ограничивает частоту создания поколений; несколько изменений за этот период объединяются в загрузку
последнего обнаруженного состояния.
Если пользовательскому ресурсу не нужен стандартный порядок подготовки, TResourceBase также
предоставляет защищённые методы MaterializeFileProvider() и MaterializeFileProviders() для
скачивания именованных файлов из уже доставленной целевой ревизии.
Чтение данных
Lock() возвращает TFileResourceAccessor<TData>, который фиксирует одну версию данных. Уже
выданный объект не меняется при переключении ресурса.
Получайте один объект доступа на одну итерацию RunIteration или на один вызов пакетной функции и
освобождайте его до следующей итерации. Долгоживущий объект удерживает старую версию и мешает
завершить переключение. Если ожидание превышает file_snapshot_rollout_warning_period, воркер
публикует ошибку /file_snapshot_activation; принудительно отозвать выданный объект нельзя.
Изменяемый файловый ресурс нельзя читать из пользовательской логики Swift-вычисления. При
повторном исполнении эпохи такое вычисление может увидеть другую версию файла и нарушить
детерминированность. Используйте материализуемое преобразование, например
TTransformComputation, либо провайдер, который гарантированно не меняется в течение всего запуска.
Закрепление версии
Провайдер может определить динамические параметры. Например,
TYTDirectoryLastFileProvider позволяет временно выбрать конкретную таблицу-ревизию вместо
последней:
dynamic_spec = {
resources = {
geobase = {
file_providers = {
release = {
parameters = {
pinned_file_name = "000001";
};
};
};
};
};
};
После изменения параметров контроллер заново обнаруживает версии. До успешного обнаружения
полного нового набора воркеры продолжают использовать прежний снимок.
Предобработка скачанных файлов
Любой файловый провайдер можно дополнить статической командой предобработки. Например, BLOB-таблица
может содержать один архив, который нужно распаковать до вызова Initialize():
file_providers = {
model = {
file_provider_class_name = "NYT::NFlow::TYTFileProvider";
parameters = {
path = "<cluster=primary>//path/to/model-archive";
};
postprocess_command = """
/usr/bin/tar -xf "$YT_FLOW_RESOURCE_PATH/model.tar" \
-C "$YT_FLOW_POSTPROCESSING_PATH"
""";
postprocess_timeout = "5m";
};
};
Flow запускает postprocess_command как /bin/bash -e -o pipefail -c <command>. По умолчанию
команда ограничена одной минутой. Ей передаётся только следующее окружение, без переменных процесса
воркера:
YT_FLOW_RESOURCE_PATH— каталог с неизменяемым скачанным деревом;YT_FLOW_POSTPROCESSING_PATH— новый пустой каталог для результата;PATH=/usr/bin:/bin,LANG=C,LC_ALL=CиTZ=UTC.
Рабочий каталог совпадает с YT_FLOW_POSTPROCESSING_PATH. Команда должна завершиться с кодом ноль
и оставить в каталоге результата только обычные файлы и каталоги; ссылки и специальные файлы
отклоняются. Команда должна синхронно дождаться всех дочерних процессов и не должна
демонизироваться. При таймауте Flow завершает всю группу процессов. Flow непрерывно вычитывает
stdout и stderr и сохраняет только последние 16 КиБ каждого потока, поэтому вывод команды не может
неограниченно увеличивать потребление памяти воркера.
Это произвольный shell-код с правами джобы воркера, а не дополнительная песочница. Команда должна
писать только в каталог результата, не должна обращаться к изменяемым внешним данным и для одной
ревизии и одной строки команды обязана получать одинаковый результат. Используйте абсолютные пути
к стабильным или версионированным исполняемым файлам. Если реализация вспомогательной программы по
тому же пути изменилась, измените видимую строку postprocess_command, например аргумент версии.
Успешный результат кешируется атомарно по идентификатору ревизии провайдера и точным байтам команды.
Попадание в кеш не повторяет ни скачивание, ни предобработку. Изменение команды инвалидирует этот
результат, но не отдельно закешированное скачанное дерево, поэтому пока оно остаётся в кеше, Flow
повторяет только предобработку. Изменение только postprocess_timeout уже готовый результат не
инвалидирует.
Ошибка команды не завершает воркер: незавершённый результат удаляется, а ресурс повторяет попытку
через file_provider_update_retry_period. На первом запуске зависимые вычисления ждут успешной
подготовки. При обновлении продолжает обслуживаться предыдущий исправный снимок. Причина доступна в
ошибке ресурса /file_update: там есть фаза, код выхода или сигнал, digest команды и ограниченные
хвосты stdout/stderr. Повторяющиеся command not found, ошибки формата входа, таймауты и падения
вспомогательной программы требуют исправить образ, данные, лимит времени или команду. Ошибки тома
и вместимости кеша дополнительно публикуются под /file_storage.
Готовые файловые провайдеры
Локальный неизменяемый файл
TLocalFileProvider предназначен прежде всего для тестов и окружений, где один абсолютный путь
указывает на одинаковый файл на всех воркерах:
file_providers = {
file = {
file_provider_class_name = "NYT::NFlow::TLocalFileProvider";
parameters = {
path = "/absolute/path/visible/to/every/worker/data.bin";
};
};
};
Файл считается неизменяемым. Не заменяйте его содержимое на месте: публикуйте новую версию по
новому пути.
Файлы в BLOB-таблице YTsaurus
TYTFileProvider материализует все файлы из одной статической сортированной BLOB-таблицы. В path
можно указать саму таблицу или link на неё. Таблица должна иметь в точности следующую строгую
схему с уникальными ключами:
<strict=%true;unique_keys=%true>[
{name="filename";type="string";sort_order="ascending";};
{name="part_index";type="int64";sort_order="ascending";};
{name="data";type="string";};
]
Каждое значение filename становится обычным файлом в корне материализованного провайдера:
file_providers = {
model = {
file_provider_class_name = "NYT::NFlow::TYTFileProvider";
parameters = {
path = "<cluster=primary>//path/to/current-model-files";
};
};
};
Части каждого файла начинаются с нуля и идут без пропусков. Имя файла должно быть одним обычным
компонентом пути.
Провайдер также поддерживает обычный файл YTsaurus: его содержимое материализуется в файл
data в корне провайдера.
Контроллер отслеживает идентификатор объекта и content_revision по заданному пути. Изменение
любого из них инвалидирует кеш и запускает reload ресурса на воркерах. При скачивании воркер берёт
snapshot-lock объекта по выбранному пути и до чтения данных сверяет его идентификатор и ревизию
содержимого с локатором. При несовпадении скачивание завершается ошибкой, и ресурс ждёт обнаружения
новой версии. Подходящую версию воркер читает в той же транзакции. Переключение link на другую
таблицу также вызывает reload, даже если числовая ревизия содержимого совпадает.
Воркеры переключаются на последнее обнаруженное поколение постепенно. Публикуйте новые версии по
отдельным неизменяемым путям, чтобы предыдущие версии оставались доступны для скачивания.
Последняя BLOB-таблица в каталоге YTsaurus
TYTDirectoryLastFileProvider воспринимает каждого непосредственного ребёнка каталога как отдельную
ревизию полного набора файлов. Провайдер выбирает ребёнка с лексикографически наибольшим именем и
материализует все файлы из выбранной BLOB-таблицы:
file_providers = {
release = {
file_provider_class_name = "NYT::NFlow::TYTDirectoryLastFileProvider";
parameters = {
path = "<cluster=primary>//path/to/releases";
};
};
};
Добавляйте новые неизменяемые таблицы под именами в виде сортируемой строки
времени, например 2026-08-31T07:00:00Z и 2026-08-31T08:00:00Z.
Динамический параметр pinned_file_name позволяет выбрать точное имя дочерней таблицы.
Дисковый кеш воркера
Каждому воркеру, который загружает файловые ресурсы, нужен worker.file_storage:
worker = {
file_storage = {
path = "/absolute/dedicated/file-resource-cache";
soft_size_limit = 1073741824;
hard_size_limit = 1342177280;
cleanup_period = "5m";
};
};
path — точный корень кеша одного процесса. Flow не добавляет к нему идентификатор пайплайна или
воркера. Процесс держит блокировку <path>/.lock и не запускается, если тот же каталог уже
используется. Поэтому тесты и окружения с общей файловой системой обязаны выдавать отдельный путь
каждому одновременно работающему воркеру.
Успешно материализованные версии переживают пересоздание ресурса и могут быть повторно использованы
после перезапуска воркера, если сохранились том и путь. Необходимые для предобработки системные
инструменты, например tar, должны присутствовать в окружении воркера.
Кеш удаляет по LRU только версии, которые сейчас не используются ресурсами. soft_size_limit —
целевой размер после очистки, hard_size_limit — граница приёма новых данных. Это не физическая
квота тома: оставляйте место на служебные данные и загрузки с заранее неизвестным размером. Если
закреплённые данные занимают больше половины жёсткого лимита, компонент публикует предупреждение.
Состояние подготовки файлов видно в статусе ресурса. Ошибка обнаружения указывает провайдер.
Ошибки скачивания, инициализации и проверки указывают ресурс, снимок и ревизии входящих в него
провайдеров. Распределение состояний снимков и отдельных ревизий публикуется метриками
/resource_controller/file_snapshot_instance_count и
/resource_controller/file_provider_revision_instance_count. Состояние кеша и ошибки нехватки
места публикуются под /file_storage.
Мониторинг файловых ревизий
На дашборде контроллера есть четыре графика для каждого ресурса и именованного файлового провайдера:
- Resource/File/Active — целевая активная ревизия контроллера: версия источника в
display_version, идентификатор кеша вrevision_id. - Resource/File/Preparing — ревизия-кандидат, ожидающая проверки.
- Resource/File/Instances — ревизии, о которых отчитались воркеры, по состояниям.
Наличие активной цели у контроллера не означает, что все воркеры уже переключились на неё. - Resource/File/Age —
Now - Timestampв секундах для целей Active и Preparing.
Если время неизвестно, серия отсутствует, а не показывает нулевой возраст.
Провайдер может заполнить необязательный TFileProviderRevision::Timestamp из метаданных источника.
Нельзя подставлять время обнаружения или скачивания. Время сохраняется при восстановлении состояния
контроллера и откате и не влияет на идентичность дискового кеша. Изменение только этого поля не
создаёт новый снимок: существующий снимок сохраняет время, записанное при его создании.
YT File и YT Directory Last используют modification_time ноды, зафиксированной snapshot-блокировкой.
Это время включает изменения метаданных, а не только содержимого. Оно не описывает время обучающих
данных модели. Local File не имеет достоверного времени публикации и оставляет поле пустым.