skip to content

You are moving an application that leans heavily on multi-key Redis commands - set intersections, MSET batches, and Lua scripts touching several keys - onto Redis Cluster. How do you decide, per access path, between co-locating keys with hash tags, restructuring the data, or moving the work into the application?

level: principalimportance: should knowfreq 32%

answer

  1. inventory access paths, don't chase CROSSSLOT errors
  2. test 1: real invariant? test 2: bounded and numerous?
  3. four moves: narrow tag / one key / split+pipeline / compute in app
  4. cross-shard = idempotent protocol, not a transaction
  5. measure slot and node skew before cutover

basics

~20 s

Classify each access path by whether it truly needs atomicity. Paths that do get a narrow hash tag; paths that merely batch get split and pipelined per node; recurring groups get collapsed into one key. Reject any co-location group that is unbounded or uneven.

solid answer

~50 s

Inventory the access paths first, then apply two tests per path. **Test 1 - does it need atomicity or a consistent snapshot?** If yes, it must run on one node: hash-tag the keys or fold them into one structure. If no (most reads, most batch writes of independent values), split it client-side and pipeline per node - a few parallel round trips, no skew. **Test 2 - is the co-location group bounded and numerous?** Per-order or per-session groups are safe: millions of small tags spread evenly. Per-tenant or global tags are not, because a slot is indivisible and cannot be rebalanced away. Where both tests bite, restructure: a hash read with `HMGET`, a precomputed set maintained at write time instead of `SINTERSTORE`, a stream per entity. Anything genuinely cross-shard becomes an application protocol - idempotent steps, reconciliation - not a fake transaction. Budget the migration by access path, and add slot-skew monitoring before cutover.

code

text · 8 lines
text
# A. invariant across two keys -> narrow tag, script is legal
EVAL "..." 2 "{order:77}:total" "{order:77}:items"

# B. fields always read together -> one key, no tag needed
HMGET user:1 name email plan

# C. independent cache reads -> split, grouped per node and pipelined
GET a  |  GET b  |  GET c      (3 keys, 2 nodes, 2 parallel round trips)

go deeper

for a junior

Recognise that multi-key commands need co-located keys, and that hash tags or a single hash are the usual fixes.

for a middle

Apply the two tests explicitly - does the path need atomicity, and is the group bounded - and know that client-side splitting is the default for non-atomic batches.

for a senior

Sequence the migration by risk, plan dual-write and backfill for renames, and define partial-failure handling for split writes.

for a principal

Treat co-location as a scarce, permanent partitioning commitment; set the policy for which invariants earn it, specify application-level protocols for cross-shard work, and consider isolating outlier tenants entirely.

## Start from an inventory, not from the errors The failure mode of these migrations is reactive: run the test suite, see `CROSSSLOT`, add a hash tag, repeat. That converges on the broadest possible tag and an unshardable cluster. Instead, enumerate the access paths that touch more than one key - grep for `MGET`, `MSET`, `*STORE`, `RENAME`, `PFMERGE`, `BITOP`, multi-key `DEL`, `EVAL` with numkeys > 1, and `MULTI` blocks - and treat each one as a design decision with a written rationale. ## Test 1: is atomicity or snapshot consistency actually required? Be strict here, because this test decides whether co-location is mandatory. - **Genuinely atomic:** invariants across keys - decrement a balance and append to a ledger; move an item between two structures; check-and-set across a pair. These must run on one node, so co-location or a single-key structure is the only option. - **Only batched:** reading twenty independent cache entries, writing a batch of unrelated values. The application would tolerate a mixed-version read, because these values are not related by an invariant. Split client-side. - **Believed atomic but isn't:** many `MGET`s that people call "consistent" were never consistent anyway - Redis offers no snapshot isolation across separate commands, and the values were written by independent commands at different times. A useful forcing question: *what invariant would be violated if these two keys were observed at different instants?* If nobody can name one, atomicity is not required. ## Test 2: is the co-location group bounded and numerous? Co-location has a hard ceiling: a slot cannot be split, so the group must always fit comfortably on one node and must not carry a disproportionate share of traffic. Score each candidate group: - **Cardinality** - how many distinct groups will exist? Millions of orders spread across 16384 slots; a dozen tenants do not. - **Size distribution** - is it roughly uniform, or a power law? Power-law groups (tenants, customers, regions) will produce a whale. - **Growth** - is the group bounded by design (an order has a fixed set of keys) or open-ended (an audit log grows forever)? Only groups scoring well on all three deserve a hash tag. If a path needs atomicity but its natural group fails these tests, that is a signal to change the data model, not to accept the skew. ## The four moves **A. Narrow hash tag.** `{order:77}:items`, `{order:77}:total`. Legal scripts and transactions, balanced slots. The default answer for genuine invariants. **B. Collapse into a single key.** Fields that always travel together become a hash (`HMGET`, `HINCRBY`); an entity's event list becomes one stream; a small set becomes one set. One key is one slot by construction, atomic operations on it are free, and no tag convention must be policed. This is the most under-used option and usually the best. **C. Split and pipeline in the client.** For non-atomic batches: group keys by owning node from the client's slot map and pipeline each group. Cost is a handful of parallel round trips instead of one; benefit is perfect balance. Verify what your client does automatically - some split `MGET`/`MSET`/`DEL` for you - and make sure your code tolerates partial failure for split writes (idempotency, retry, or a compensating pass). **D. Move the computation up.** Some multi-key commands exist because Redis was the easiest place to compute something. A `SINTERSTORE` of two large sets can often be replaced by maintaining the intersection incrementally at write time, or by doing the intersection in the service where the two members are already in memory. This trades a hot node for application CPU, and often reduces total work. ## Cross-shard operations that survive Some workflows genuinely span shards - a transfer between two accounts that must not be co-located, a fan-out write. Redis will not give you a distributed transaction, and pretending otherwise is the real risk. Design an explicit protocol: idempotent operations keyed by a request id, a durable intent record in the system of record, retries, and a reconciliation job. Decide the failure semantics you are willing to expose (at-least-once with dedupe is the usual answer) and write them down. ## Cost, sequencing and evidence - **Sequence by risk:** paths needing atomicity first (they constrain key names), then batched reads, then computation moves. - **Key renames are the expensive part.** Plan dual-write plus backfill for any path whose keys change name, and finish the migration before cutting over routing. - **Instrument before you migrate:** per-master memory and key counts, per-slot key counts for your largest tags, and per-node command rates. Skew is trivial to see and ruinous to discover late. - **Consider isolation for outliers.** If two customers are 100x the rest, a dedicated cluster for them is a legitimate, cheaper answer than distorting the whole keyspace to accommodate them. ## What a good answer sounds like Not "use hash tags", and not "avoid multi-key commands". It is: *co-location is a scarce resource; spend it only where an invariant demands it and only on groups that are numerous and bounded; convert everything else into single-key structures, client-side fan-out, or application logic; and treat true cross-shard work as a protocol problem with explicit failure semantics.*

  • How do you tell a team that a required cross-shard atomic operation simply cannot be provided?
    Show that a transaction or script executes on one node against local slots only, so cross-shard atomicity has no implementation in Redis Cluster. Then offer the two real options: co-locate the keys if the group is bounded and numerous, or specify an application protocol with idempotent steps, a durable intent record and reconciliation, and agree explicitly on at-least-once semantics with deduplication.
  • What evidence would make you choose a dedicated cluster for a handful of customers rather than a cleverer key scheme?
    A power-law size or traffic distribution where the top few groups are one to two orders of magnitude above the median, so any co-location scheme that keeps their operations atomic also creates an unshardable slot. Isolation gives them their own capacity, keeps the shared cluster balanced, and avoids distorting the key naming of every other access path.

Co-location is like reserved seating on a train: worth it for a family that must sit together, ruinous if you reserve a whole carriage for one company - the seats still exist, but nobody else can use them and you cannot split the carriage.

saying these in an interview costs you the question

  • Resolving CROSSSLOT errors reactively until every key shares one broad tag
  • Claiming a multi-key read was atomic when the values were written independently anyway
  • Assuming a Lua script can be made to span shards with enough cleverness
  • Ignoring partial-failure semantics after splitting a multi-key write client-side
  • Deferring key-name decisions until after the data is large, when renames require dual-write and backfill

context