Scanning a Keyspace That Won't Hold Still: How SCAN Iterates a Cluster
A single-key GET is easy to reason about: hash the key, find its owner, ask that one node. SCAN is the opposite kind of operation. It has no key — it asks a question about the whole keyspace: give me every key in this range, a page at a time. And in a partitioned, replicated store that keyspace isn't one thing in one place. It's split across 256 segments spread over however many nodes you run, each node holding only the slices it currently owns, and the whole set is being written to while you iterate. There is no global sorted index to walk. So where does the ordering come from, what exactly does a resume token point at, and what happens to your scan when a key is written — or a partition moves — halfway through it?
This post follows a SCAN down through the three layers that actually serve it: a single segment, a single node, and the client that stitches the nodes together. Each layer solves one piece of the problem, and the honest limits of the operation fall out of how those pieces fit.
The request: a range and a page size, not a key
A scan is described by two small objects. A KeyRange says which keys — a start, an end, or a prefix, any of them optional:
public sealed class KeyRange
{
public string? Start { get; set; }
public string? End { get; set; }
public string? Prefix { get; set; }
}
And ScanOptions says how to deliver them — how big a page, whether to ship values or just keys, an optional hard limit, and a resume token to continue a previous call:
public sealed class ScanOptions
{
public int PageSize { get; set; } = 100;
public bool IncludeValues { get; set; } = true;
public bool IncludeMetadata{ get; set; } = false;
public string? ResumeToken { get; set; }
public int? LimitCount { get; set; }
}
Nothing here names a node or a segment. That's the point: the caller describes what they want, and the three layers below work out where it lives and in what order to hand it back.
Layer one: a segment scans itself, lock-free and in order
The smallest unit that can answer part of a scan is a segment — Zaris's unit of partitioned storage, backed by a ConcurrentDictionary of live records. A dictionary has no order, so the first thing a segment scan needs is a sorted view of its keys. Building that on every page would be ruinous, so each segment keeps a cached, immutable sorted-key snapshot and rebuilds it only when the key set changes:
// No lock held during scan — the snapshot is immutable once built.
var sortedKeys = GetOrBuildSortedKeys();
int startIdx = BinarySearchKey(sortedKeys, startKey, StringComparer.Ordinal);
if (startIdx < 0) startIdx = ~startIdx;
The array is assigned once and never mutated; a concurrent rebuild just produces a new array and swaps the reference. That means the scan walk itself takes no lock at all — it binary-searches to the start key and walks forward in ordinal order, reading live records out of the concurrent dictionary as it goes. Reads and writes to the segment never block the scan, and the scan never blocks them.
Two kinds of key are silently passed over as it walks:
if (!_data.TryGetValue(key, out var record)) continue; // vanished since snapshot
if (record.IsDeleted || ExpirationManager.IsExpired(record, now)) continue; // tombstone / TTL
A deleted key is a tombstone — a real record marked IsDeleted, because a delete has to be a record, not an absence, for replication to carry it. The scan filters those out, along with keys whose TTL has elapsed but whose sweeper hasn't fired yet. So the sorted-key snapshot is a superset of what you'll see; expiry and deletion are resolved against the live record at read time, not frozen into the snapshot.
The walk stops at PageSize and returns a ResumeToken. Here is the first honest detail: at the segment level, the token is essentially an offset — a count of how many matching entries were already returned (ReturnedSoFar). A resumed segment scan re-binary-searches to the same start key and skips that many matches before emitting. It's a position-by-count, not a position-by-key seek.
cursor.ReturnedSoFar += taken;
bool hasMore = results.Count >= pageSize && cursor.ReturnedSoFar < totalLimit;
That choice matters for what the operation promises, and we'll come back to it.
Layer two: one node merges its segments into a single ordered stream
A node owns many segments — not all 256, just the ones mapped to the partitions it currently serves. A scan at the node level has to present all of them as one ordered result, because the caller asked about a range, not about segment 42. This is a classic merge problem: each segment can produce its keys in sorted order, and you want the sorted union.
So the first page of a node scan does exactly that. It gathers a fetcher for every locally-owned segment and runs a k-way merge across them with a priority queue, comparing keys ordinally:
var heap = new PriorityQueue<(KvResult<ZarisEntry> Item, int Source), KvResult<ZarisEntry>>(comparison);
// seed with the head of each segment, then repeatedly pop the smallest
// key and pull the next from that segment's enumerator
The merged, ordered result is materialized into a buffer held by a RouterQueryCoordinator — a short-lived, server-side session registered under a freshly minted QueryId:
coord = new RouterQueryCoordinator();
await coord.InitScanAsync(fetchers, ct); // drains + merges every owned segment
var id = _coordManager.Register(coord);
cursor = new QueryCursor { QueryId = id, ReturnedSoFar = 0 };
Every subsequent page is just a slice of that buffer, and the node's resume token carries the QueryId so the next call finds the same session:
if (!string.IsNullOrEmpty(options?.ResumeToken))
{
cursor = QueryCursor.Deserialize(options.ResumeToken);
if (!_coordManager.TryGet(cursor.QueryId, out coord))
throw new InvalidOperationException($"Search session {cursor.QueryId} not found (expired or cleaned).");
}
This is the second honest detail, and it's a real trade-off. A node-level scan is stateful: paging is stable and the per-node ordering is exact, but the price is that the node materializes its entire owned, merged keyspace into a buffer when the scan begins, and holds a session for the life of the iteration. Scans are an administrative and bulk-iteration tool, not a constant-memory streaming cursor over billions of keys — and a resume token is only good against the node and session that minted it. Let that session expire, or try to resume it somewhere else, and you get a loud "session not found" rather than a quietly wrong answer.
Layer three: the client fans out and interleaves the owners
One node answers for the segments it owns. A cluster scan has to cover every owner, so the client-side ScanClusterClient opens a scan against each node in parallel and interleaves their pages as they arrive, rather than draining one node before starting the next:
var enumerators = clients
.Select(c => c.Scan.ScanAsync(range, options, ct).GetAsyncEnumerator(ct))
.ToList();
// advance all owners concurrently; yield whichever produces next
var finishedTask = await Task.WhenAny(pending.Select(p => p.Task));
That parallelism is why LimitCount needs care. Each node treats the limit as a per-node ceiling; if the client didn't also enforce it, a LimitCount of 1,000 against four owners could return up to 4,000 rows once merged. So the cluster layer applies the limit as a hard total cap and stops yielding once it's met — while still draining every owner's in-flight page to completion, because disposing an async enumerator mid-move isn't allowed and the per-node results are already bounded anyway:
if (emitted < limit)
{
yield return enumerator.Current;
emitted++;
}
// keep draining siblings even past the cap; do not yield-break mid-move
The consequence to be clear about: within any one node's contribution, keys arrive in order; across nodes they're interleaved by arrival, so a full cluster scan is not globally sorted unless you ask for it explicitly. (The related Search path does carry SortFields and merges with a comparator for exactly that reason.) If you need every key in the cluster in ordinal order, sort the assembled result; SCAN itself guarantees per-owner order and completeness, not a global sort.
What a SCAN does — and does not — promise
Put the three layers together and the operation's real contract is visible, without any of it being hidden:
- It is not a point-in-time snapshot of the cluster. Each node materializes its own merged buffer when its scan begins, and those moments don't coincide. Writes that land on a node before its buffer is built are included; writes after it are not. There is no cluster-wide freeze.
- Deleted and expired keys never surface. Tombstones and TTL-elapsed records are filtered against the live record as the segment is walked, so a scan reflects the logical, not the physical, contents.
- Resume is position-by-count inside a session. Because the per-segment cursor is an offset, a key inserted or removed before the current position between two pages can nudge the boundary — a scan concurrent with heavy mutation of the range being scanned can, in principle, repeat or skip a key around the seam. A scan of a quiescent range is exact. This is the same family of guarantee Redis
SCANgives — a best-effort full iteration, not a snapshot isolation level — and it's the right trade for an unordered, concurrently-written store. - The keyspace can move under you. The client fans out to the owners that exist when the scan starts. Partitions can migrate or fail over mid-iteration; a scan is a best-effort sweep of a live cluster, which is why it's the tool for bulk export, reindexing, and administrative queries rather than for transactional reads. For read-your-writes on a single key, route to the owner instead.
None of these are bugs to be papered over; they're the honest shape of iterating an unordered, sharded, mutating keyspace without stopping the world. The design makes the cheap, common case fast — a lock-free sorted snapshot and a per-node merge — and is explicit about the edges.
A note on the RESP front-end
Redis clients reach the same machinery through KEYS and SCAN on the RESP front-end, which map onto the native scan. One caveat worth stating plainly: the RESP SCAN today is single-shot — it drains the native scan and replies with cursor 0 (iteration complete) in one round trip, rather than threading Redis's incremental cursor through each call. It's honest about completeness, not yet incremental. The native client is the one that pages with a resume token.
One operation, three honest layers
A GET touches one owner and one key; a SCAN touches every owner and asks each to produce its slice of an ordered whole. The work splits cleanly: a segment turns an unordered dictionary into a sorted, lock-free walk; a node merges its segments into one ordered, stably-paged stream behind a session cursor; the client fans out across owners, interleaves their pages, and enforces the one cap that fan-out would otherwise break. Each layer is simple on its own, and the operation's guarantees — per-owner order, completeness over a quiescent range, no snapshot isolation, no global freeze — are exactly what those three simple layers can honestly deliver together.