---
metadata:
  - name: generator
    content: Diplodoc Platform v5.50.6
alternate:
  - https://ytsaurus.tech/docs/en/flow/concepts/key_visitor.md
  - https://ytsaurus.tech/docs/ru/flow/concepts/key_visitor.md
---
> **Documentation Index:** Fetch the complete configuration index at https://ytsaurus.tech/docs/en/llms.txt

<!-- source: en/_includes/flow/concepts/key_visitor.md -->
# Key Visitor Streams in YTsaurus Flow

## Why you need key-visitor streams {#why-key-visitors}

Stateful computations accumulate an [internal state](https://ytsaurus.tech/docs/en/flow/concepts/glossary.md#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](https://ytsaurus.tech/docs/en/flow/concepts/timers.md)) 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 {#how-it-works}

### Pass lifecycle {#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 {#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 {#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 {#supported-computations}

| Process-function adapter | Visit stream support |
|---|---|
| `TProcessFunctionComputation` | ✓ |
| `TProcessFunctionSwiftMapComputation` | ✓ (only for state handling: emitting to output from `ProcessVisit` is forbidden, see [Swift](swift.md#swift-map)) |
| `TProcessFunctionSourceComputation` | ✗ |
| `TProcessFunctionTransformOrderedSourceComputation` | ✗ |

## Configuration {#configuration}

To make a computation accept a visit stream, you must fill the `key_visitor_streams` field in its static spec:

```yson
"tester" = {
    "computation_class_name" = "NYT::NFlow::TProcessFunctionComputation";
    "processing_function" = "NYT::NFlow::NMyProject::TMyFunction";
    "group_by_schema" = [...];
    "key_visitor_streams" = {
        "visit_iter" = {};
    };
    ...
};
```

### Schema requirements {#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 {#static-params}

- `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](https://ytsaurus.tech/docs/en/flow/concepts/glossary.md#state) and visitor-driven external-state joiners (see [Static table joiner](#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](#pass-lifecycle)). When unspecified, all input and source streams of the computation. Only meaningful with `finite=%true` (see [Dynamic spec parameters](#dynamic-params)). For when the list needs narrowing, see [`streams_dependency`](#deps).

### Dynamic spec parameters {#dynamic-params}

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](#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 {#key-visitor-only}

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](#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](#schema-requirements)) is mandatory for such a computation—the partition ranges are built from it.

Set `finite=%false` in the [dynamic parameters](#dynamic-params) 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.

{% note 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.

{% endnote %}

### Static table joiner {#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](https://ytsaurus.tech/docs/en/flow/cpp/state.md#static-table-key-visitor-joiner) section.

### `streams_dependency` {#deps}

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:

```yson
"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](https://ytsaurus.tech/docs/en/flow/python/testing.md) and the pipeline must reach [`completed`](https://ytsaurus.tech/docs/en/flow/concepts/glossary.md#start-stop-pause-pipeline), that is, in [integration tests](https://ytsaurus.tech/docs/en/flow/python/testing.md). There, keep in `upstream_streams` (see [Static spec parameters](#static-params)) only the input streams that don’t depend on the visitor:

```yson
"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 {#observability}

### Visit EventTimestamp {#event-timestamp}

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 {#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 {#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.
<!-- endsource: en/_includes/flow/concepts/key_visitor.md -->