Persistence and sharding
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 hot layer (RAM). The authoritative simulation of online players and loaded zones runs in the memory of the server process at the tick rate (positions, combat, NPC spawns). Fast, ephemeral.
- The cold layer (a durable database). A relational database (SQL Server / MySQL / Postgres): characters, inventory as foreign keys to items, guilds, mail. Slow, reliable. It holds everyone — including offline characters, who cannot be kept in RAM.
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 seconds. A node crash between flushes loses up to seconds of progress — the player logs back in "in the past" (the classic "server rollback"). The expected loss with a uniform cache is:
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: 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 per "game second", while its real capacity per wall second is . The load is . While everything is fine; as it approaches 1 the action queue grows catastrophically, per queueing theory:
That used to be the "death lag" of big fights: , the queue never drains, the node dies. TiDi (2011) introduces a dilation factor : simulation time runs at speed relative to real time. Then per wall second the node only has to do of work. We set:
The effective load — the queue is bounded again, and everything is processed in order and fairly. The floor (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:
A load six to ten times over unity → 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".
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.
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.
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).
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 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 , and mean time in the system follows from Little's law . At there are ~9 requests in the system; at — ~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 , we divide the rate at which actions "mature" for processing, holding .
Why a floor at all, rather than "as much as it takes"
Without a floor, under extreme load would go to ~0 — the fight would never end, and the node would keep accumulating jitter anyway. A hard floor of 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.
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 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.
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.EVE slows time down instead of throwing hardware at an overloaded system — why can't you just add CPUs?
WoW has orders of magnitude more players than EVE — so why is EVE considered technically harder?
If one shard gives you a single world and a shared history, why doesn't everyone do it?
Item dupes are "hackers", surely. Why do you call it an engineering bug?
Why have a database at all — why not keep the whole world in RAM, since that is faster?
- CCP Games, the dev blogs "Introducing Time Dilation (TiDi)" (2011) and "Time Dilation — How's That Going?" — the primary sources on the slowdown mechanic.
- "The Bloodbath of B-R5RB" (eveonline.com) — a breakdown of the largest PvP battle; the payoff of a single shard.
- Leslie Lamport, "Paxos Made Simple" (2001) — state consistency between database servers.
- Tim Sweeney, "Network Architecture for Massively Multiplayer Games" (2001) — the move from P2P to authoritative client-server.
- Warcraft Wiki: "Sharding", "Layering", "Cross-realm zones" — the vocabulary of subdivision in theme-park MMOs.
- Module 4 (
04-online-worlds-1997-2005.md), sections "Database Design for Persistent Worlds" + "EVE Online's Distributed Architecture".