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

  1. Background fill. Each partition runs a background loop. It reads the state page by page using KeyStates::List and passes the keys to an internal visit buffer. The speed is regulated by a throttler configured so that one full pass takes the specified Period.
  2. Emit. Ready TVisit messages are delivered by the engine via GetNextBatch and reach the process function in ProcessVisit.
  3. 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 table key_visitor_states. That’s why a worker restart or partition rebalance doesn’t cause a re-scan.
  4. 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 takes Period. 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.
  5. Final pass. With the default finite=%true, when every followed stream is Completed (by default, all input and source streams—see upstream_streams), the visitor finishes after a final pass and its partitions become Completed. With the default full_final_pass=%true, that pass starts after the inputs finish and guarantees at least one complete scan. Setting full_final_pass=%false marks the current pass final and can finish sooner without that guarantee. With finite=%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_schema must start with a uint64 column. 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 via clear() in the visit handler. If only external_names is specified (without names), only the listed external state is scanned; internal states aren’t scanned. If neither names nor external_names is specified, all internal states are scanned, but external state and joiners aren’t; a joiner is scanned only if explicitly listed in external_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 with finite=%true (see Dynamic spec parameters). For when the list needs narrowing, see streams_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 single KeyStates::List call (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 and GetNextBatch.
  • finite: whether the visitor follows the computation inputs and sources to completion. Default %true performs a final pass after they finish, then completes the partitions. Set %false for 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 %false resumes scanning only before a pass has been marked final; once partitions are Completed, they must be recreated by repartitioning.
  • full_final_pass: whether the final pass must be a complete scan begun after the inputs finish. Default %true guarantees 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—about 1 / background_fill_period per second per source per partition (with the default of 500 ms, that’s about 2 reads/s), independent of period (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_iteration is too small. If a single key has more internal-state names than max_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_stall via StatusProfiler (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 increase max_scan_rows_per_iteration via Reconfigure.
  • 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.