Key Visitor Streams in YTsaurus Flow
Why you need key-visitor streams
Stateful computations accumulate an internal state per key while processing incoming messages. You often need to periodically scan this state—even without an incoming message for a specific key. Common tasks include:
- TTL / eviction: iterate over all keys and delete outdated records.
- Periodic aggregates: emit the current value for each key once a day.
- Forced re-evaluation: trigger a recalculation for all keys in the background.
Standard timers (see Timers) aren’t suitable here. You must register them for each key in advance, but the set of keys can grow without an explicit “new key appeared” event. For example, keys may arrive only via state stores, not through a message stream.
A key-visitor stream solves this problem. A background task in the worker periodically scans the entire state of a partition and emits a TVisit message for each key into a special internal stream. The process function subscribes to this stream via ProcessVisit (C++) or on_visit / OnVisit / onVisit (Python / Go / Java) and decides how to handle its state, just like it would for a regular incoming message.
How it works
Pass lifecycle
- Background fill. Each partition runs a background loop. It reads the state page by page using
KeyStates::Listand passes the keys to an internal visit buffer. The speed is regulated by a throttler configured so that one full pass takes the specifiedPeriod. - Emit. Ready
TVisitmessages are delivered by the engine viaGetNextBatchand reach the process function inProcessVisit. - Coverage. After visits for a key range are delivered to the consumer, the range is marked as Committed in
TKeyVisitorStore. The coverage is persisted to the system tablekey_visitor_states. That’s why a worker restart or partition rebalance doesn’t cause a re-scan. - End of pass. When the coverage is complete, the background loop immediately calls
StartNewPass. The pace is set by the throttler, so the next pass still takesPeriod. Rotation is atomic: a single Sync-transaction deletes the previous pass’s rows and seeds the first interval of the new one. If a crash happens midway, it rolls back, and the coverage is preserved. - Final pass. With the default
finite=%true, when every followed stream is Completed (by default, all input and source streams—seeupstream_streams), the visitor finishes after a final pass and its partitions becomeCompleted. With the defaultfull_final_pass=%true, that pass starts after the inputs finish and guarantees at least one complete scan. Settingfull_final_pass=%falsemarks the current pass final and can finish sooner without that guarantee. Withfinite=%false, there is no final pass and the visitor keeps scanning.
Partitioning
The state is partitioned by the uint64 hash of the first column in group_by_schema (usually farm_hash(key)). Each partition uses its own TKeyVisitor to scan its hash range.
Inside a partition, the range is split into a statically defined number of buckets (bucket_count). Buckets are scanned in round-robin order. This smooths the load and keeps the “unfinished coverage” evenly distributed across the partition. That way, if a worker crashes, it doesn’t discard progress in one narrow area.
Guarantees
| Property | Guarantee |
|---|---|
| Each key per period | You get exactly one visit per key (no duplicates or omissions). |
| Worker restart | Coverage is preserved; rotation is atomic, and a crash keeps the old pass. |
| Partition rebalance | The new worker sees the committed coverage via key_visitor_states. |
| Completion of followed streams | With default finite=%true, the visitor completes after the inputs and sources selected by upstream_streams finish. Default full_final_pass=%true guarantees a full pass started after completion; %false may finish during the current pass. |
finite=%false |
No final pass: the visitor scans while the pipeline runs. |
Period |
Best-effort. Under throttler load or slow KeyStates reads, the achieved period grows. See observability below. |
| Scan order | Within a bucket, keys are sorted; between buckets, the order is round-robin. Not event-time. |
Which computations support visit streams
| Process-function adapter | Visit stream support |
|---|---|
TProcessFunctionComputation |
✓ |
TProcessFunctionSwiftMapComputation |
✓ (only for state handling: emitting to output from ProcessVisit is forbidden, see Swift) |
TProcessFunctionSourceComputation |
✗ |
TProcessFunctionTransformOrderedSourceComputation |
✗ |
Configuration
To make a computation accept a visit stream, you must fill the key_visitor_streams field in its static spec:
"tester" = {
"computation_class_name" = "NYT::NFlow::TProcessFunctionComputation";
"processing_function" = "NYT::NFlow::NMyProject::TMyFunction";
"group_by_schema" = [...];
"key_visitor_streams" = {
"visit_iter" = {};
};
...
};
Schema requirements
group_by_schemamust start with auint64column. This column is used as the hash to split the partition into buckets. The check runs at spec submit time.
Static spec parameters
names: a list of internal-state names to count keys from. If unspecified or empty, all keys from all internal states of the computation are used.external_names: a list of external-state managers and visitor-driven external-state joiners (see Static table joiner) for the computation, whose tables are used to count keys. This lets the visitor scan the external state (including that of companion computations) and evict outdated records viaclear()in the visit handler. If onlyexternal_namesis specified (withoutnames), only the listed external state is scanned; internal states aren’t scanned. If neithernamesnorexternal_namesis specified, all internal states are scanned, but external state and joiners aren’t; a joiner is scanned only if explicitly listed inexternal_names, and in no more than one visit stream of the computation (checked at spec submit time).upstream_streams: which of the computation’s input and source streams the visitor follows, that is, whose completion it waits for before the final pass (see Pass lifecycle). When unspecified, all input and source streams of the computation. Only meaningful withfinite=%true(see Dynamic spec parameters). For when the list needs narrowing, seestreams_dependency.
Dynamic spec parameters
In dynamic_spec.computations.<id>.key_visitor_streams.<name>:
period: the target duration for one full pass. Default is 1 day.max_scan_rows_per_iteration: the limit for a singleKeyStates::Listcall (in rows, not keys). It must be strictly greater than the maximum number of internal-state names per key; otherwise, the scan stalls (see Diagnostics).buffer_row_limit: the maximum size of the internal buffer of ready visits between the background fill andGetNextBatch.finite: whether the visitor follows the computation inputs and sources to completion. Default%trueperforms a final pass after they finish, then completes the partitions. Set%falsefor a periodic scanner that must run throughout the pipeline lifetime, especially a computation with no input or source streams. This dynamic setting can be switched on a running pipeline to request completion. Switching back to%falseresumes scanning only before a pass has been marked final; once partitions areCompleted, they must be recreated by repartitioning.full_final_pass: whether the final pass must be a complete scan begun after the inputs finish. Default%trueguarantees at least one such scan. With%false, the current pass can be marked final and stop sooner, without that guarantee.background_fill_period: the pause between iterations of the background loop in idle state. Each such iteration performs one read (KeyStates::List) per scanned source. So this parameter sets the base frequency of the visitor’s read requests—about1 / background_fill_periodper second per source per partition (with the default of 500 ms, that’s about 2 reads/s), independent ofperiod(which only affects the width of the hash slice for each read). Under load (cap-hit, bucket/pass change), iterations are rescheduled immediately, and the frequency can be higher.
Computation with only a visit stream
A key-visitor stream can be the only source of work for a computation: input_stream_ids is empty, there are no source_streams, and work comes only from scanning the external state (external_names). This is the auditor pattern: periodically re-evaluate keys in an external table without a message stream. Such a computation is partitioned by uint64 hash ranges in the same way as an input-driven one (see Partitioning); the number of partitions is set by min_partition_count / max_partition_count / desired_partition_count in the dynamic spec.
The requirement for group_by_schema (first column must be uint64, see Schema requirements) is mandatory for such a computation—the partition ranges are built from it.
Set finite=%false in the dynamic parameters for this input-free scanner. With default %true, there are no inputs to wait for, so its first pass is final and the partitions complete after one full scan.
Warning
One external table, one writer. A visitor computation modifies the external state. Don’t modify the same external table from multiple computations: concurrent writes from different computations aren’t serialized and lead to races and lost updates. At spec validation time, YT path ownership is checked: if two computations declare writes (via their external-state managers) to the same (cluster, path), the spec submit fails with the claimed for writing error.
Static table joiner
The scan can read not only the computation’s state but also an external static table—via the external-state joiner TStaticTableKeyVisitorJoiner, listed in external_names. The table must be strictly sorted by the computation’s group_by_schema: the prefix of its key columns must match that schema in names and types. The joiner reads the table sequentially, in the same key order as the state scan, and passes the table row to ProcessVisit as a read-only state for the visit key. Table keys take part in the scan on an equal footing with state keys: a visit arrives even for a key that isn’t yet in the computation’s own state.
This is the basis of the reconciliation pattern: a periodic scan aligns the computation’s own state with the external table—keys present in the table are updated, and keys missing from it are deleted. Requirements for the table, behavior when the source is unavailable, and a code example are in the TStaticTableKeyVisitorJoiner section.
streams_dependency
By default, the visit stream is not part of the computation’s streams_dependency: visitors are usually for internal cleanup and don’t produce output, and a stuck visitor shouldn’t block the completion of output streams. The “last pass” signal is delivered to the visitor locally by the worker (SetUpstreamCompleted)—this doesn’t require an edge in the graph.
If the process function emits messages to output from ProcessVisit, you must explicitly list the visit stream as the parent of that output in streams_dependency. Example:
"streams_dependency" = {
"visits" = ["keys"; "visit_iter"];
};
Here, visits is the output stream, keys is the input, and visit_iter is the key-visitor stream. For a computation with only a visit stream (no input), there’s a single parent: "visits" = ["visit_iter"].
If that output comes back to the same computation as an input stream, you get a cyclic topology: the cycle visit_iter → requests → … → responses closes on the computation’s own input stream. By default, the visitor follows that input, but it only completes once the visitor stops scanning the state.
In a production pipeline this doesn’t matter: sources are infinite, input streams never complete, and no value of upstream_streams starts a final pass. It matters where all pipeline sources are finite and the pipeline must reach completed, that is, in integration tests. There, keep in upstream_streams (see Static spec parameters) only the input streams that don’t depend on the visitor:
"input_stream_ids" = ["keys"; "responses"];
"output_stream_ids" = ["requests"];
"key_visitor_streams" = {
"visit_iter" = {
"upstream_streams" = ["keys"];
};
};
"streams_dependency" = {
"requests" = ["keys"; "visit_iter"];
};
Here, responses is produced by the visitor via requests, so the visitor follows keys only.
Observability
Visit EventTimestamp
Each emitted TVisit carries an EventTimestamp equal to the expected scan time for that key. The formula is:
scheduleLag = max(0, elapsed − ScannedFraction · Period)
EventTimestamp = now − scheduleLag
where elapsed is the time since the start of the current pass (taken as the minimum PassStartedAt among already scanned intervals; it’s persisted in key_visitor_states, so it survives restarts and rebalances).
If the scan runs on schedule, scheduleLag = 0, and EventTimestamp ≈ SystemTimestamp. If you fall behind (for example, the throttler limits throughput due to a slow backend), EventTimestamp lies in the past, and the visit stream’s watermark lags by exactly the amount of the delay. This gives a direct signal to the downstream consumer and shows up as event-lag in standard flow metrics.
Scheduled visit rate
In the visit stream's inflight metrics, new_count_per_sec and offered_count_per_sec both estimate the rate required to visit all matching keys within period: estimated distinct key count divided by the period in seconds. This is not the actual buffer fill rate. ready_count reports buffered visits, and processed_count_per_sec reports committed processing.
The estimate uses the density of distinct matching keys in recently scanned hash ranges: smoothed key count divided by smoothed hash coverage, multiplied by the partition's hash span. Both counters use the same observation timestamps and a 30-second EMA window. Before the window matures, the ratio of cumulative counts and coverage provides an initial estimate. Empty reads contribute coverage with zero keys and can lower the estimated density; an estimated zero does not mean that the visitor is Empty.
The counters are not reset between passes. With no new reads, including while the buffer is full, the estimate stays unchanged. No population or EMA state is persisted: after restart or repartitioning, the job estimates density from new local observations, independently of the restored scan cursor. Without positive hash coverage, the rate is absent; this includes ranges with identical hash values at both boundaries, even after a complete pass. Small samples, uneven hash distribution, and capped reads within one hash can cause transient bias. Changes in unscanned regions are not immediately visible. Once the finite visitor becomes Empty, its New and Offered rates are zero.
Diagnostics
max_scan_rows_per_iterationis too small. If a single key has more internal-state names thanmax_scan_rows_per_iteration, the scan can’t progress: the limit is hit mid-key, all its rows are discarded, and no progress is made. The computation reports the error/key_visitor/<stream>/scan_cap_stallviaStatusProfiler(Key visitor stalled: a single key has more than max_scan_rows_per_iteration = N rows). The error clears automatically once the read progresses. The solution is to increasemax_scan_rows_per_iterationviaReconfigure.- State backend is unavailable. If reading the state during a scan fails, the worker doesn’t crash: an error is reported to
/key_visitor/<stream>/background_fill, and the read is automatically retried. A transient YT/RPC error resolves on retry; a persistent issue (schema mismatch, missing table) remains visible in the status and doesn’t break the pipeline.