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):

  1. Читает TUrlMessage с полями Host и Url
  2. Сохраняет URL в стейт хоста
  3. Регистрирует хост и URL в TUrlDownloader
  4. Ставит таймер для периодической проверки результатов через GetNextHostCheck
  5. Применяет лимит на размер стейта через EnforceLimit

При срабатывании таймера (ProcessTimer):

  1. Восстанавливает хост из стейта
  2. Извлекает обработанные URL из TUrlDownloader
  3. Удаляет обработанные URL из стейта
  4. Генерирует TProcessedUrlMessage для каждого обработанного URL
  5. Если очередь пуста — отменяет регистрацию хоста и сбрасывает стейт; иначе — ставит следующий таймер

Типы сообщений

  • 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, будут потеряны. Значение лимита следует выбирать с учетом допустимых потерь.

Структура пайплайна

  1. Входная очередь → поток urls (TUrlMessage)
  2. Поток 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") — входные URL
  • RegisterStream<TProcessedUrlMessage>("processed_urls") — обработанные URL

Исходный код

TUrlDownloader

TLimitedUrlDownloadFunction

См. также

Предыдущая
Следующая