URL Downloader в YTsaurus Flow (C++)
Пайплайн демонстрирует выполнение внешних вызовов (подобных HTTP-запросам) из компьютейшенов с управлением скоростью загрузки (троттлинг), шардированием по хостам и управлением стейтом.
Компоненты пайплайна
TUrlDownloader
Вспомогательный класс (наследник TRefCounted), управляющий очередями загрузки для каждого хоста. Обеспечивает:
- Регистрацию хостов (
RegisterHost) — для каждого хоста создается асинхронный исполнитель, который последовательно обрабатывает URL-адреса из очереди с искусственной задержкой (эмуляция троттлинга) - Добавление URL (
RegisterUrl) — добавляет URL в очередь соответствующего хоста - Извлечение результатов (
ExtractProcessedUrls) — возвращает обработанные URL с результатами - Отмену хоста (
UnregisterHost) — останавливает обработку и очищает очередь
Ключевой паттерн — использование AsyncVia(GetCurrentInvoker()) для запуска фоновой обработки в рамках сериализованного инвокера, что позволяет безопасно работать с общим состоянием без блокировок.
TLimitedUrlDownloadFunction
Основная process function реализует IProcessFunction и запускается адаптером TProcessFunctionComputation. Она координирует загрузку URL с помощью TUrlDownloader.
При обработке входного сообщения (ProcessMessage):
- Читает
TUrlMessageс полямиHostиUrl - Сохраняет URL в стейт хоста
- Регистрирует хост и URL в
TUrlDownloader - Ставит таймер для периодической проверки результатов через
GetNextHostCheck - Применяет лимит на размер стейта через
EnforceLimit
При срабатывании таймера (ProcessTimer):
- Восстанавливает хост из стейта
- Извлекает обработанные URL из
TUrlDownloader - Удаляет обработанные URL из стейта
- Генерирует
TProcessedUrlMessageдля каждого обработанного URL - Если очередь пуста — отменяет регистрацию хоста и сбрасывает стейт; иначе — ставит следующий таймер
Типы сообщений
- TUrlMessage — наследник
TYsonMessage. Содержит поляHostиUrl. - TProcessedUrlMessage — наследник
TYsonMessage. Содержит поляHost,UrlиData(результат обработки).
Ключевые паттерны
Internal YsonState
Для хранения очереди URL по каждому хосту используется TMutableStateKeyClient<TLimitedHostState>. Стейт TLimitedHostState содержит:
Host— имя хостаUrls— очередь URL (std::deque<std::string>), ожидающих обработки
Таймеры для периодической проверки
Метод GetNextHostCheck вычисляет время следующей проверки хоста. Время вычисляется с учетом:
CheckHostPeriod— период проверки (по умолчанию 5 секунд)- Хэш имени хоста — для равномерного распределения проверок разных хостов по времени
Динамические параметры
TDynamicLimitedUrlDownloadParameters, переданный через динамическое поле processing_function_parameters, позволяет менять параметры без перезапуска пайплайна:
CheckHostPeriod— период проверки хостов (по умолчанию 5 секунд, должен быть больше 1 секунды)PersistLimit— максимальное количество URL, сохраняемых в стейте для одного хоста (по умолчанию 1000)
PersistLimit
EnforceLimit ограничивает размер очереди URL в стейте. Если количество URL превышает PersistLimit, старые URL удаляются. Это необходимо, чтобы стейт не превышал ограничения на размер строки в динамической таблице.
Важно
При ребалансировке партиций URL, не попавшие в лимит PersistLimit, будут потеряны. Значение лимита следует выбирать с учетом допустимых потерь.
Структура пайплайна
- Входная очередь → поток
urls(TUrlMessage) - Поток
urls→ TLimitedUrlDownloadFunction (с таймерами и стейтом) → потокprocessed_urls(TProcessedUrlMessage)
В статической спеке для url_downloader указаны computation_class_name = "NYT::NFlow::TProcessFunctionComputation" и processing_function = "NYT::NFlow::NExample::TLimitedUrlDownloadFunction". Динамические настройки функции передаются через dynamic_spec/computations/url_downloader/processing_function_parameters.
Функция main
В main регистрируются два потока, а TLimitedUrlDownloadFunction регистрируется через YT_FLOW_DEFINE_PROCESS_FUNCTION:
RegisterStream<TUrlMessage>("urls")— входные URLRegisterStream<TProcessedUrlMessage>("processed_urls")— обработанные URL