URL Downloader in YTsaurus Flow (C++)
The pipeline shows how you run external calls (similar to HTTP requests) from computations. You control the download speed (throttling), shard by host, and manage the state.
Pipeline components
TUrlDownloader
This helper class (a descendant of TRefCounted) manages download queues for each host. It lets you:
- Register hosts (
RegisterHost) — for each host, it creates an asynchronous executor that processes URLs from the queue one after another with an artificial delay (to emulate throttling). - Add URLs (
RegisterUrl) — it adds a URL to the corresponding host’s queue. - Extract results (
ExtractProcessedUrls) — it returns processed URLs with their results. - Unregister a host (
UnregisterHost) — it stops processing and clears the queue.
The key pattern is using AsyncVia(GetCurrentInvoker()) to run background processing within a serialized invoker. This lets you safely work with shared state without locks.
TLimitedUrlDownloadFunction
This is the main process function. It implements IProcessFunction, is hosted by TProcessFunctionComputation, and coordinates URL downloads using TUrlDownloader.
When you process an input message (ProcessMessage):
- You read a
TUrlMessagewith theHostandUrlfields. - You save the URL in the host’s state.
- You register the host and URL in
TUrlDownloader. - You set a timer to periodically check results via
GetNextHostCheck. - You apply a limit to the state size using
EnforceLimit.
When the timer fires (ProcessTimer):
- You restore the host from the state.
- You extract processed URLs from
TUrlDownloader. - You remove processed URLs from the state.
- You generate a
TProcessedUrlMessagefor each processed URL. - If the queue is empty, you unregister the host and reset the state; otherwise, you set the next timer.
Message types
- TUrlMessage — a descendant of
TYsonMessage. It contains theHostandUrlfields. - TProcessedUrlMessage — a descendant of
TYsonMessage. It contains theHost,Url, andDatafields (the processing result).
Key patterns
Internal YsonState
You use TMutableStateKeyClient<TLimitedHostState> to store the URL queue for each host. The TLimitedHostState state contains:
Host— the host name.Urls— the URL queue (std::deque<std::string>) waiting to be processed.
Timers for periodic checks
The GetNextHostCheck method calculates the time for the next host check. The time is calculated based on:
CheckHostPeriod— the check period (default is 5 seconds).- The host name hash — to evenly distribute checks for different hosts over time.
Dynamic parameters
TDynamicLimitedUrlDownloadParameters, passed through the dynamic processing_function_parameters field, lets you change parameters without restarting the pipeline:
CheckHostPeriod— the host check period (default is 5 seconds and must be greater than 1 second).PersistLimit— the maximum number of URLs stored in the state for a single host (default is 1000).
PersistLimit
EnforceLimit limits the size of the URL queue in the state. If the number of URLs exceeds PersistLimit, the older URLs are removed. This keeps the state within the row size limits for a dynamic table.
Warning
During partition rebalancing, URLs that don’t fit within the PersistLimit will be lost. Choose the limit value considering the acceptable loss.
Pipeline structure
- Input queue → the
urlsstream (TUrlMessage). - The
urlsstream → TLimitedUrlDownloadFunction (with timers and state) → theprocessed_urlsstream (TProcessedUrlMessage).
In the static spec, url_downloader uses computation_class_name = "NYT::NFlow::TProcessFunctionComputation" and processing_function = "NYT::NFlow::NExample::TLimitedUrlDownloadFunction". Function settings are passed through dynamic_spec/computations/url_downloader/processing_function_parameters.
main function
In main, you register two streams; TLimitedUrlDownloadFunction is registered via YT_FLOW_DEFINE_PROCESS_FUNCTION:
RegisterStream<TUrlMessage>("urls")— input URLs.RegisterStream<TProcessedUrlMessage>("processed_urls")— processed URLs.