A larger screen, mouse and keyboard make designing and running simulations much easier. Open this page on your laptop or desktop to continue.
How backend and distributed-system ideas connect. Pick one to see how it works.
Capacity and performance approximation using thought experiments and common performance numbers to evaluate architectural designs.
Temporary in-memory data store intercepting frequent read queries before they hit the database.
Map servers and keys onto one hash ring so that adding or removing a server moves only about k/n of the keys.
A geographically dispersed network of servers used to deliver static content near users to improve load times.
Horizontally partitioning a large database into smaller, independent shards based on a shard key.
Scaling capacity by adding more compute instances into a resource pool behind a load balancer.
Relative latency profiles of computer operations across memory, network, and disk that guide system architecture.
Evenly distributes incoming network traffic among web servers defined in a load-balanced set.
A durable in-memory buffer supporting asynchronous communication that decouples producers and consumers.
A mechanism that limits the frequency of client requests over a defined time window to prevent DoS attacks, reduce third-party API costs, and protect backend servers from overload.
Decoupling user session data from web server memory into external shared storage to allow seamless autoscaling.
A structured 4-step collaborative process (Scope, High-Level Design, Deep Dive, Wrap Up) to navigate open-ended system design interviews.
A rate limiting algorithm where tokens accumulate in a bucket at a fixed refill rate up to a capacity, allowing traffic bursts while bounding long-term average throughput.
Scaling a system by adding more compute, memory, or disk capacity to an existing single server machine.
Keeping requests that read and write the same data at the same time from overwriting each other's changes.
How fresh reads are: strong (never out of date), weak (may miss the latest write) or eventual (replicas agree given enough time).
Separating write operations on a primary database from read queries distributed across replica nodes.
A non-relational store of unique keys and opaque values with put and get; spread over many servers it needs partitioning, replication and tunable consistency.
N copies, W acknowledgements per write, R answers per read: W + R > N guarantees strong consistency, small W or R answers faster.
Evaluating traditional relational RDBMS databases with SQL joins against specialized non-relational stores.
Writes go to a commit log and a memory cache that is flushed to sorted SSTables; reads try memory, then a bloom filter to pick SSTables.
An operation that has the same effect whether it is applied once or several times, so it can be retried safely.
A distributed system keeps at most two of consistency, availability and partition tolerance; since partitions happen, it chooses CP or AP.
Decentralized failure detection: nodes pass heartbeat counters to random peers and mark a member offline when its counter stops growing.
System operational continuity measured as uptime percentages or 'nines' defined formally by Service Level Agreements.
During temporary failures, the first healthy servers on the ring stand in for down replicas and hand the data back when they return.
Replicas compare trees of hashes, root first, to find and sync only the buckets that differ after a permanent failure.
Deploying systems across multiple geographic regions with GeoDNS routing for low latency and disaster recovery.
Writes a number with the 62 characters 0-9, a-z, A-Z; a URL shortener turns each unique ID into its short URL this way.
A 64-bit ID built from a 41-bit timestamp, a 5-bit datacenter ID, a 5-bit machine ID and a 12-bit sequence, so IDs are unique and sortable by time.
Hands out IDs that are unique across many servers; the chapter compares multi-master auto_increment, UUID, a ticket server and Twitter snowflake against one set of requirements.
Gives a long URL a short alias and redirects the alias back; reads outnumber writes 10 to 1, so redirects are served from a cache in front of the database.
[server, version] pairs on each data item that tell whether two versions descend from each other or conflict.
Each server takes many places on the hash ring, so keys spread more evenly; more virtual nodes, smaller standard deviation.