URL Downloader в YTsaurus Flow (Java)
Пайплайн группирует входящие URL по хосту, накапливает их во внутреннем стейте и обрабатывает пакетами по таймеру: через 5 секунд после постановки URL в очередь срабатывает таймер, который эмитирует результаты в выходной стрим.
Компоненты
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.