← Module 4/Persistence and sharding
RU
Module 4 · Online worlds (1997–2005)

Persistence and sharding

An authoritative server is where the truth lives. But when there are thousands of players, the truth does not fit in one process. The era gave two answers: cut the players (sharding — thousands of independent copies of the world, like WoW's realms) or cut the space (one world spread across cluster nodes, like EVE) — and pay for it by slowing time down under load.
~17 min☁ infra + databases
The gist in 30 seconds
A persistent world = authoritative state that survives a disconnect, a crash and a patch while serving thousands of mutations per second. The architecture has two layers: a hot authoritative simulator in RAM (tick by tick) plus a cold durable store (a relational database) that the world is periodically "saved" into. One process cannot carry 100k players, so there are two ways to scale. Sharding (WoW): cut the population — N independent copies of the world, each with its own database (the partition key = the realm). It scales linearly with hardware but fragments the social graph: players on different realms cannot see each other. One shard (EVE's "Tranquility"): cut the space — one universe divided into solar systems, each assigned to a cluster node; everyone interacts across the whole world, but the throughput of any one system is capped by a single node. When 2670 pilots pile into one system (as in B-R5RB, 2014), the node chokes → Time Dilation: slow the simulation down to 10% of real time so that every action can be processed in order. That is backpressure: trading speed for correctness. The price of sharding is a fractured world; the price of one shard is heroic engineering and bullet time in a big fight.

The mechanism: where the truth physically sits

The networking model of an MMO is authoritative client-server (see the taxonomy): the server is the single source of truth, the client only predicts and draws. This lesson is about the next question: on which hardware that truth lives when concurrency is in the tens of thousands, characters in the database number in the millions, and every piece of state has to survive the client shutting down.

Two layers: a hot simulator and a cold store

World state splits in two:

The world is periodically saved: hot state is serialized into cold state (on logout, on a timer, on important transactions). Between saves the durable copy lags behind the live one. Hence the defining class of bugs of the era: rollbacks and item dupes on a crash, when save and transaction boundaries are sloppy.

Saving, the rollback window and dupes

Say the world is flushed to the database every T seconds. A node crash between flushes loses up to T seconds of progress — the player logs back in "in the past" (the classic "server rollback"). The expected loss with a uniform cache is:

E[Lloss] ≈ T2 ; T=300 s ⇒ on average 2.5 min lost

A dupe is more dangerous than a rollback, because it destroys the economy. Handing over an item is two mutations: debit A, credit B. If they are not atomic and they straddle a flush boundary when the crash happens, then on recovery the item may end up both with A (rolled back to the state "before the trade") and with B (the credit made it to disk) — it has doubled. The cure is not "anti-cheat" but making trade an ACID transaction over the durable store, instead of a lazy write-behind of two independent character blobs.

One process is not enough — two ways to cut

A node hits its CPU/memory/database limits long before 100k players. The scaling fork:

Sharding (WoW): cut the population world A DB A world B DB B world C DB C realms are independent · key = realm A and C do NOT see each other One shard (EVE): cut the space node node node node node node systems of the universe → nodes one DB · one world One system overloaded: 2670 pilots on one node the action queue grows → lag Time Dilation: stretch the tick clock until it all fits 100% 50% 10% (the floor) — bullet time simulation speed s falls as load rises; every action is processed in order and fairly, the world just runs in slow motion. The cost ≠ correctness.

Sharding: cutting the population

Sharding is horizontal partitioning of players: the world is copied into N independent instances (realms / servers), each with its own authoritative state and its own database. This is literally database sharding where the partition key = the realm. WoW is built exactly that way: on the character select screen you choose a realm — that is, a database partition. Upsides: linear scaling (need more room — spin up another realm), one realm failing does not take the rest down, and a small consistency "surface". There is one downside, but a heavy one: the social graph is fractured — friends on different realms cannot trade or raid together; moving a character is a paid migration of a row between databases. Modern mitigations: connected realms, cross-realm zones (CRZ — the reverse, merging underpopulated realms so a zone feels alive), FFXIV's data centers with "World Visit".

One shard: cutting the space

EVE swam against the current: one universe for everyone (the "Tranquility" server), partitioned not by population but by space. The universe is carved into ~8000 solar systems; at any moment each system "belongs" to one node (process) in a large cluster. Players interact across the whole world — one market, one history — but the throughput of an individual system is capped by one node: a fight in a system cannot be "parallelized" across two CPUs, because every ship in it interacts with every other (this is not embarrassingly parallel). So the only lever under overload is the rate at which time passes.

Time Dilation: backpressure by slowing time

EVE runs the server tick at 1 Hz (the simulation advances by 1 second per tick). Say a node has to process work W per "game second", while its real capacity per wall second is C. The load is L=W/C. While L<1 everything is fine; as it approaches 1 the action queue grows catastrophically, per queueing theory:

Qavg ≈ L1−L (L→1⇒Qavg→∞)

That used to be the "death lag" of big fights: L≥1, the queue never drains, the node dies. TiDi (2011) introduces a dilation factor s: simulation time runs at speed s relative to real time. Then per wall second the node only has to do s·W of work. We set:

s= max(0.1, min(1, CW ))

The effective load Leff=s·W/C≤1 — the queue is bounded again, and everything is processed in order and fairly. The floor s=0.1 (10%) guarantees that the fight will eventually end: 1 second of simulation per 10 seconds of wall clock. TiDi does not touch EVE server time (industry, skill training, structure timers run on the real clock) — only "physics in space" gets stretched.

A worked example — why B-R5RB ran for 21 hours

In "The Bloodbath of B-R5RB" (January 2014) up to 2670 pilots were in one system at once (7548 involved in total), over 21 hours, with 576 capitals destroyed including 75 Titans, ~11 trillion ISK (≈$300–330k in real money). Suppose a node can cleanly handle the actions of ~250 ships in real time (illustrative, not a published CCP number). Then:

L= 2670250 ≈10.7 ⇒ s=max(0.1,1/10.7) =0.1

A load six to ten times over unity → s pinned to the 10% floor. The fight literally runs at one tenth speed — hence the bullet time in which a Titan's salvo takes minutes to land, and 21 hours of real time contain only a few hours of "in-game" battle. Without TiDi the node would simply have died; with TiDi the battle played out slowly but deterministically and fairly. That is the payoff of one shard: B-R5RB is a real historical event, not "an episode on server #47", because there is only one universe.

🕹 Games to play — and what to notice

You can see the persistence architecture with the naked eye — on the server select screen (or its absence). Play up the ladder: from an explicit "pick a shard" to "there is no shard, there is one world that knows how to slow down".

World of Warcraft sharding by realm

The reference implementation of "cut the population". When you create a character you pick a realm — which is picking a database partition. Inside a realm the zone is subdivided further (sharding a zone under a rush, layering a whole realm at the Classic 2019 launch, later removed).

🎮 Play: create characters on two different realms — notice they are in separate universes: they cannot see each other, cannot trade, and share only the account. Then look at the price of "Character Transfer" in the Blizzard store — you are paying for migrating a row between databases. That is the cost of sharding, written out in dollars.

RuneScape / Lineage / most MMOs an explicit "world select"

The same model, even more plainly: a list of "worlds"/"servers" with a population indicator right in the lobby. Every world is an independent shard; the popular ones are packed, the empty ones sit idle.

🎮 Play: in RuneScape open the world list — dozens of numbered shards with a player counter. Log into a low-population world and a packed one: the economy, Grand Exchange prices and crowding are noticeably different, because these are different worlds, glued together only by a shared account service.

Final Fantasy XIV megaserver · data centers · World Visit

The modern compromise. Logical "worlds" are grouped into data centers; inside a world, zones are instanced (sharding under the hood), but you can visit between worlds of the same DC (World Visit). The social graph is far less fractured than in classic WoW.

🎮 Play: in FFXIV do a World Visit to a neighboring world in your data center — notice that instanced zones (Limsa, the markets) split into "duplicates" under load, and that you have physically moved onto a different world process. A hybrid: it cuts both the population (worlds) and the space (instances).

EVE Online one shard, Tranquility + TiDi

The opposite pole: you never chose a server at all. One universe, one market, one history for everybody. The price is Time Dilation in a big fight.

🎮 Play: log into EVE and fly to a system with a known large fleet (or watch a siege stream). Find the TiDi indicator (the clock icon / speed percentage) — you will see it drop toward 10% as the system fills up, and everything around you slide into slow motion. Compare with a quiet system (100%). You are literally watching backpressure: the server slows time down so as not to lose a single action.

Deep end · engineering: hot/cold, write-behind, and where dupes come fromskippable

The heart of persistence is the boundary between the hot authoritative state in RAM and the cold durable store. How you cross it determines both performance and an entire class of economy bugs.

Write-through vs write-behind

  • Write-through: every significant mutation is committed to the database immediately. Reliable, but the database becomes the bottleneck at 10k+ transactions/s — you cannot push every step a character takes through SQL.
  • Write-behind (typical for MMOs): mutations accumulate in RAM, and the database is flushed in batches on a timer or at logout. Fast, but it opens a rollback window T and demands care about write ordering.

So you split by importance: position/HP go write-behind (losing them is not scary), while money and items go write-through inside a transaction (losing or doubling them is a catastrophe).

Anatomy of a dupe

Handing over an item is two writes: A.inventory -= item and B.inventory += item. If that is not one transaction but two independent write-behind blobs, and the node dies at the moment when B += item has been persisted but A -= item has not (or recovery replays the log in a "convenient" order) — the item ends up with both. Nearly every real dupe exploit in UO/Diablo/WoW is of that class: not "a hacker" but broken atomicity of a distributed transaction at the seam between the hot and cold layers, often provoked by a crash/disconnect at just the right moment.

The cure

  • Trade/auction — an ACID transaction over the durable store (a single writer per item partition, or two-phase commit if the characters live in different database shards).
  • Idempotency and monotonic operation IDs — replaying the log after a crash must not credit twice.
  • Transfer of ownership through an atomic change of a foreign key, not "delete from one / create for the other".
Deep end · hosting: cluster topology, reinforced nodes and the price of the choiceskippable

Where you "cut the space", placing load onto hardware is exactly gameplay infrastructure.

A node = the owner of spatial cells

In a single-shard architecture the cluster is a pool of processes (in EVE they are historically called SOL nodes), and a dispatcher spreads solar systems across nodes. Usually many quiet systems live on one node; a hot system can get a node to itself. Moving a system between nodes (or merging them) is a non-trivial operation: all the authoritative state has to be handed over without "double ownership".

A reinforced node — the scheduled siege

EVE lets you request node reinforcement ahead of time: if alliances have announced a battle over a system, CCP moves it onto a dedicated, more powerful node before it starts. That is manual "horizontal pre-scaling" for a predictable peak — a luxury available precisely because there is one world and events in it are public.

The price of the two models

  • Sharding scales almost linearly and cheaply: a new realm = a box of hardware plus a database instance, operationally simple. Which is why mass theme-park MMOs (WoW, thousands of realms) choose it.
  • One shard demands heroic engineering (TiDi, dynamic node placement, handoff handling) and tolerates bullet time. It only pays off if a single world is the point of the game (a sandbox, emergent politics, one market).

The business consequence of both: when population declines, sharded games do server merges (welding empty realms together — the same pain as CRZ), while single-shard games are immune to this by construction.

Deep end · theory: why the queue explodes as L→1 and why the floor is exactly 10%skippable

TiDi is admission control on top of a queueing system. A node processing player actions is modeled as a queue: arrivals at rate λ, service at rate μ, load ρ=λ/μ.

The nonlinearity at unity

For M/M/1 the mean number in the system is N=ρ/(1−ρ), and mean time in the system follows from Little's law N=λ·W. At ρ=0.9 there are ~9 requests in the system; at ρ=0.99 — ~99. That is why the lag of big fights is not linear in the number of players: while there is headroom it is barely noticeable, and at the boundary it falls off a cliff. TiDi scales the effective arrival rate: by slowing the simulation clock by a factor of s, we divide the rate at which actions "mature" for processing, holding ρeff≤1.

Why a floor at all, rather than "as much as it takes"

Without a floor, under extreme load s would go to ~0 — the fight would never end, and the node would keep accumulating jitter anyway. A hard floor of s=0.1 guarantees progress: even in the worst case one second of battle takes ten seconds of wall clock. It is a deliberate trade: under absurd load a node may still start falling behind even at the floor (actions then get buffered), but the game stays coherent and fair — nobody "teleports", the order of actions is preserved. Backpressure instead of data loss.

Analogy
A restaurant chain. Sharding — open 50 identical branches: each with its own kitchen and register, the queue splits, scaling is trivial — but your friends at "branch B" physically cannot sit at your table in "branch A", and transferring a reservation between branches is a separate procedure. One shard — one giant restaurant for everybody: you can run into anyone, one menu, one history of the place. But when a crowd piles into one room, the waiters cannot keep up — and the place switches on "slow time": it serves everyone in turn, fairly, just slower, so that nobody is lost and no orders get mixed up. Nobody leaves without food; the evening just drags on.
Why it matters
Persistence is where game development runs headlong into a pure distributed system: where the authoritative truth lives, how it survives a failure, and what to do when load exceeds one node. The choice "cut the population or cut the space" is not about graphics but about what kind of world the game can be at all: fragmented and easy to scale, versus unified and expensive. And TiDi is a rare honest answer to "what do we do when the hardware physically cannot keep up": do not lie to the client, do not lose actions, but slow time down and keep correctness. The same choice — shard for simplicity or hold a single state at the cost of heroics — faces anyone building systems outside games too.
🔁 Beyond games — where this transfers
The lesson is two canonical distributed-systems moves plus backpressure, shown on a live example.

ML / AI: these are the two axes of training parallelism, word for word. Data parallelism = sharding: every worker holds a full copy of the model and its own slice of the batch (like realms — independent copies of the world), with synchronization via all-reduce/parameter server = "connected realms". Model / tensor parallelism = one shard cut along space: one logical model sawn across GPUs (like EVE's universe across nodes), and activations at layer boundaries are the inter-system handoffs (expensive cross-device exchange). When a device or a collective saturates, gradient accumulation and microbatching play the role of TiDi: we lower the effective step rate to stay correct within the memory/interconnect limit, instead of dropping updates. And the checkpoint cadence is exactly the "save every T" trade-off: checkpoint rarely and you lose up to T of progress when a node dies. Sharded KV/feature stores for online inference are the same partition-by-key.

Databases: sharding = horizontal partitioning by key (consistent hashing, the hot-shard problem, resharding); single-writer-per-partition versus two-phase commit across shards is literally "the dupe in a cross-shard trade". The WoW↔EVE dilemma is an AP flavor (independent partitions) against a CP flavor (a single state with graceful degradation).

Distributed systems / infra: cell-based architecture (AWS cells = shards, failure isolation), sticky sessions, stateful service placement = assigning systems to nodes. Backpressure / load shedding (slow the producer down instead of losing messages) is exactly TiDi: slow the input, keep correctness.

Business: server merges as concurrency declines = merging underloaded regions/tenants; the cost of fragmenting your base (users or data) is always higher than it looks at the start.

The principle: when load exceeds one node, choose deliberately — cut into independent copies (simple, but it fragments) or hold a single state (powerful, but it needs backpressure at the peak). And never lose the truth silently: better to slow down than to lie.

🔧 Run it and poke at it — on your home machine
What to play is above (🕹). This part is for people who want to touch the boundary between the hot and cold layers with their hands.
🔧 Poke at it (debug) ~40 min, local Postgres
Spin up a local Postgres/SQLite. Model a "trade" as two updates: UPDATE a SET item=NULL and UPDATE b SET item=:x. First run them without a transaction and "kill" the process (Ctrl-C / kill) in between — then look at the state: the item is either lost or doubled depending on the order. Then wrap it in BEGIN … COMMIT and crash again — the invariant "the item is in exactly one place" holds. That is the anatomy of a dupe, live.
🧪 Test it (with QA eyes) ~15 min, EVE / a stream
Find a system under load in EVE (or on a siege stream) and debug TiDi by eye: at what pilot count does the indicator leave 100%? Does it reach the 10% floor? Do skill/structure timers stop while it does (they should not — that is EVE time, which is not dilated)? Then a thought-experiment QA pass on handoff: a pilot flies from system A (node 1) to B (node 2) — which invariant stops the ship from existing in both at once during the transfer?
Checklist: reproduced a dupe without a transaction and fixed it with an ACID wrapper; watched TiDi fall to 10% and verified that server time is not dilated; stated the ownership-transfer invariant for a node handoff.
Connections
foundation
Netcode taxonomy — authoritative client-server: what we put on the wire and who the source of truth is. This lesson is where that truth physically sits and how it survives scale.
foundation
Client prediction — the authoritative server up close; persistence is the same authority made durable and spread across a cluster.
contrast
Rollback netcode — P2P with no authority, a tiny ephemeral state, a save every frame so you can roll back to move forward; here it is authority, an enormous durable state, and saving so the world survives. The same "frequent serialization" for opposite goals.
next
Virtual economies — EVE's single shard is exactly what makes one market and a measurable economy possible; sharding, by contrast, splits the economy into N independent ones.
Questions worth asking
EVE slows time down instead of throwing hardware at an overloaded system — why can't you just add CPUs?
Because one solar system is one node, and a fight inside it is not embarrassingly parallel: every ship is shooting at every other, applying modules, computing ranges — the state is densely coupled. Splitting that system across two CPUs means shipping nearly all the state between them every tick (the cross-talk eats the gain) and solving consistency in real time. It is cheaper and more correct to keep the system on one node and pull the only safe lever — the rate at which time passes. "Add hardware" works between systems (they are spread over nodes), but not inside a single battle.
WoW has orders of magnitude more players than EVE — so why is EVE considered technically harder?
Because WoW never has 7000 people who can shoot at each other in one consistent state. WoW is thousands of small independent worlds (realms), plus zone sharding: the load is split by construction. EVE deliberately chose one world — and so signed up for the problem WoW sidesteps architecturally: holding a single authoritative state with thousands of interacting entities in one place. More players in total ≠ harder; harder is when they are all in one consistent simulation.
If one shard gives you a single world and a shared history, why doesn't everyone do it?
Because for a theme-park MMO, fragmenting the social graph is a feature, not a bug: realms keep the population dense and manageable and the content drip measured. A single world buys a WoW designer nothing (the content is instanced and scripted anyway) but costs TiDi-grade engineering and a design that tolerates bullet time. One shard only pays off where a single sandbox, emergent politics and one market are the whole point of the game (EVE). It is a product choice that forces an architecture choice, not the other way round.
Item dupes are "hackers", surely. Why do you call it an engineering bug?
Because the overwhelming majority of historical dupes are broken transaction atomicity at the seam between hot RAM state and the cold database, not broken cryptography. The player only provokes the right timing (a disconnect/crash/double click at the moment of transfer); the doubling itself is done by the system, replaying the log in a "convenient" order or partially persisting two independent mutations. So the fix is not "ban the cheaters" but making trade one ACID transaction with idempotency. It is exactly the distributed-transaction class of bugs, just dressed up as an inventory.
Why have a database at all — why not keep the whole world in RAM, since that is faster?
Durability and volume. Durability: processes crash, servers reboot for patches, a data center can blink — without a durable store, any such moment wipes everybody's progress. Volume: you can only keep online players in loaded zones hot; millions of offline characters with all their inventory will not fit in RAM and do not need to be there. So the layering is mandatory: the hot layer is whoever is playing right now, the cold layer is the truth about everyone. The only question is the cadence and the atomicity of moving between them — and that is where rollbacks and dupes live.
Further reading