skip to content

After assignments move, what does a caller with a stale partition map observe, and why is that worse than an error?

level: seniorimportance: must knowfreq 52%

answer

  1. the assignment goes stale, not the rule
  2. wrong node, ordinary-looking absence
  3. follow it and learn from it
  4. mid-move is per call only
  5. absence is not proof of absence

basics

~20 s

It depends on whether the node reached checks ownership. Where it does not, a misrouted read returns an ordinary absence and the caller treats live state as gone; where it does, a reply names the right node to follow and learn from.

solid answer

~50 s

A stale copy of the map means a caller sends a key to a node that no longer owns it, and the outcome is a property of the store. Where nodes do not track which partitions they own, the read comes back **absent** — indistinguishable from a key that was never stored — so a session looks logged out, a record marking work as already handled looks missing and the work runs twice, a claim looks free. Nothing logs an error. Where nodes do track ownership, the node answers that the key lives elsewhere, and the caller must both follow that answer and refresh its map, or it pays two hops on every operation forever. During a live assignment move there is a third case: the node says this particular key is in flight, which the caller should follow for that call only and must not write into its map. An error is the lucky outcome; silence is the one that corrupts behaviour.

go deeper

for a junior

Know that a caller can hold an out-of-date idea of which node owns a key, and that the result is a request sent to a node that will not have it.

for a middle

Explain the two obligations a reply that names the right node creates: follow it now, and refresh the map so the next operations go straight there.

for a senior

Show the silent case and design for it: a misroute that returns an ordinary absence makes the application treat live state as gone, with no error anywhere to trigger a retry or an alert.

for a principal

Set the standing position that absence on this tier is not proof of absence, and decide which state is allowed to depend on a lookup that can silently lie.

## What is actually stale The partition map is the rule that turns a key into a partition plus the current assignment of partitions to nodes. The rule does not go out of date; **the assignment does**, whenever a node is added, a node fails, or partitions are moved between nodes while the tier keeps serving. A caller that holds its own copy of the map is therefore holding a snapshot, and between the moment assignments change and the moment the caller learns of it, that caller resolves keys to the wrong node. What happens next is not one behaviour. It is one of three, and which one you get is a property of the store rather than of the idea. ## Outcome one: silence On stores whose nodes do not check which partitions they own — which includes a large part of this class, where partitioning is entirely a client-side convention over independent nodes — the node reached simply does what it was asked: - **A read** finds nothing under that key, because the value lives on the node that now owns the partition. It replies *absent*. The caller cannot distinguish that from an entry that was never written, or one that expired, or one that was evicted. - **A write** succeeds. It is stored on the wrong node, where no correctly routed caller will ever look, and it occupies memory until its lifetime, if it has one, runs out. This is the failure mode that matters most here, and it is worth being precise about why. The application does not see an error and therefore does not retry, alert or fall back to a system that could reproduce the state. It sees an ordinary miss and acts on it: the session is treated as expired and the user is sent to log in again; the record marking a request as already handled is treated as absent and the request is processed a second time; a quota reads as unspent; a claim reads as unheld. **Live state is quietly treated as gone**, at a rate proportional to how many keys moved. ## Outcome two: a reply that the key lives elsewhere Stores whose nodes track their own partitions can do better. The node answers that the key is not here and names where it is. That answer carries two obligations for the caller, and designs routinely honour only the first: 1. **Follow it** for the operation in hand, so the request still succeeds. 2. **Learn from it** by refreshing the map, so the next thousand operations for that partition go straight to the right node. A caller that follows without learning is correct and slow: it pays an extra hop on every affected operation indefinitely, and the symptom is a latency rise with no error rate to explain it. The rate of such replies is, incidentally, the cleanest signal that a fleet is routing on old assignments. ## Outcome three: the key that is mid-move While partitions are being moved and the tier keeps serving, a single key can be in flight — not yet fully on the new node, no longer authoritative on the old one. Stores that can express this answer differently: *this key, right now, try over there*. The distinction matters because the two answers deserve opposite treatment. - A **permanent** answer means the assignment changed: follow it and update the map. - An **in-flight** answer means the assignment has not changed yet: follow it for this call and leave the map alone. A caller that writes an in-flight answer into its map routes a whole partition to a node that does not own it yet. A caller that refuses to follow one sees absences for the keys that have already moved. ## When the caller holds no map Behind an intervening proxy, or with a directory consulted per lookup, the caller has no copy of the map and so cannot hold a stale one. That does not mean the problem is gone — it means the staleness moved. A proxy with an out-of-date assignment misroutes on the caller's behalf, producing exactly outcome one or two at the proxy, with the caller unable to see which. The advantage is real, though: there is one place to refresh instead of one per client process. ## Seeing it in production | Symptom | Likely cause | Where to look | |---|---|---| | Miss rate rises while request volume is flat, right after a topology change | callers resolving on an old assignment | miss ratio against traffic, correlated with the change | | Latency rises with no error rate | callers following the reply but not refreshing the map | rate of replies saying the key lives elsewhere | | Duplicate work, double charges, repeated notifications | absences read as truth for records with no other home | the workloads whose state exists nowhere else | | One partition's keys unreachable after a move | an in-flight answer cached as permanent | the caller's map version against the current assignment | The design conclusion follows from outcome one: on this kind of tier, **absence is not proof of absence**. Where the consequence of being wrong is severe, either the state must be verifiable against a system that can reproduce it, or the routing shape must be one where a misroute cannot be silent. ## The three branches a caller-side resolver must distinguish. Only the first updates the map; the second is per call; the third is the one that carries no information about whether the key exists ``` node = route(key, map) // map is this caller's snapshot reply = node.get(key) if reply is redirect_moved: // the assignment changed for good map = refresh_map() // learn it, or pay two hops forever reply = route(key, map).get(key) else if reply is redirect_in_flight: // this one key is mid-move reply = reply.target.get(key) // follow for this call only // do NOT write it into map if reply is absent: // where nodes do not check ownership, a misroute lands here // instead: live state read as though never stored handle_miss(key) ```

  • Why is a misrouted write often worse than a misrouted read?
    A read returns an absence and the caller usually recovers by recomputing. A write on the wrong node succeeds, so the caller believes the state is stored; correctly routed callers never see it, and it consumes memory until its lifetime, if any, runs out. The damage stays invisible until something depends on that state being there.
  • How do you keep a design safe when a caller's map may be stale?
    Assume any lookup can return an absence that is not the truth. Refresh on the first reply that says the key lives elsewhere rather than on a schedule alone, never treat absence as proof for state whose loss is expensive, and make the records that must not vanish verifiable against a system that can reproduce them.
  • Which signal tells you a fleet is routing on old assignments?
    The rate of replies saying a key lives elsewhere, where the store produces them, and otherwise the miss ratio measured against request volume around a topology change. A latency rise with a flat error rate points to callers following such replies without refreshing their copy of the map.

A delivery office working from an old address book fails in two different ways. If it simply attempts the old address, the letter comes back marked no such resident, and the sender cannot tell a person who moved from a person who never lived there — so the sender concludes the recipient does not exist. If the old address instead carries a forwarding note, the letter arrives and the sender learns the new address for next time. Same stale book, but one version produces an answer you can act on and the other produces an absence you cannot interpret.

saying these in an interview costs you the question

  • Assumes every store answers a misroute by naming the right node
  • Follows the reply but never refreshes the map
  • Treats an in-flight answer as a permanent assignment change
  • Says a misrouted write fails, so nothing can be stored wrongly
  • Reads an absence as proof the state was never written
  • Thinks a proxy in front makes stale assignments impossible