URL Downloader в YTsaurus Flow (Java)

Пайплайн группирует входящие URL по хосту, накапливает их во внутреннем стейте и обрабатывает пакетами по таймеру: через 5 секунд после постановки URL в очередь срабатывает таймер, который эмитирует результаты в выходной стрим.

Исходный код (Java)

Исходный код (Kotlin)

Компоненты

UrlDownloadFunction

Основная процессная функция, реализующая логику накопления и обработки URL. Метод onMessage добавляет URL в внутренний стейт хоста и устанавливает таймер; метод onTimer обрабатывает накопленные URL и эмитирует результаты:

HostState

Модель внутреннего стейта, сериализуемая в YSON. Хранит имя хоста и список URL, ожидающих обработки:

Регистрация компьютейшена и стримов

Компьютейшен url_downloader регистрируется аннотацией @FlowComputation на классе process-функции:

Типизированные стримы объявляются через ComputationProvider (метод getStreams()):

PipelineMain

Единственная точка входа (запускает пайплайн или обслуживает его как компаньон — по YT_FLOW_MODE):

Ключевые паттерны

  • Группировка по ключу — стейт создаётся отдельно для каждого хоста; Flow автоматически направляет сообщения с одинаковым ключом в один экземпляр компьютейшена.
  • Таймер на основе wall-clock времени — output.addTimer(System.currentTimeMillis() / 1000 + 5, 0L) запускает обработку через 5 секунд после постановки URL в очередь. Несколько вызовов addTimer с одинаковым triggerTimestamp дедуплицируются.
  • Пакетная обработка в onTimer — все накопленные URL обрабатываются разом при срабатывании таймера, что снижает число обращений к downstream-сервисам.
  • Очистка стейта — после обработки стейт удаляется через accessor.clear(), предотвращая утечку памяти.
  • YsonStateAccessor — внутренний стейт сериализуется в YSON и хранится на стороне C++ воркера; Java-объект получается через getOrDefault.

См. также

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