System Design
A practitioner's reference for designing distributed systems: requirements and constraints, scaling and partitioning, consistency models, synchronous versus event-driven communication, failure domains, backpressure, and knowing when not to distribute.
System design is the practice of deciding how work and state are distributed across components so that the whole satisfies its requirements under load and partial failure. It is mostly the study of tradeoffs: every technique that buys throughput, availability, or independence costs consistency, complexity, or operational burden somewhere else. A design is not good or bad in the abstract — it is well or badly matched to a set of constraints.
What system design is
The things a system designer actually chooses are narrower than the diagrams suggest:
- Where each piece of state lives, and who is allowed to write it.
- Which interactions are synchronous and which are asynchronous.
- How the system is partitioned, and along what key.
- What the system does when a component is slow, gone, or returning wrong answers.
- Which guarantees are promised to users, and which are deliberately not.
Note what is not on that list: technology names. Choosing between two databases matters far less than choosing the access pattern they will serve, and a design expressed only as a set of product names has not made any of the decisions above.
Start with constraints, not components
Designs go wrong early, by drawing boxes before establishing what the system must do. Four categories of constraint do almost all the work:
- Functional requirements. What operations exist, and what each must guarantee.
- Workload shape. Read-heavy or write-heavy; steady or spiky; uniformly distributed or dominated by a few hot entities; batch or interactive. These shapes point at different architectures more decisively than scale does.
- Service objectives. Acceptable latency at the tail, acceptable unavailability, acceptable staleness, acceptable data loss on failure. Until these are stated, no tradeoff can be evaluated — and "as fast and reliable as possible" is not a constraint, it is an abdication.
- Hard limits. Regulatory data location, existing systems that cannot change, team size, budget, and — most underrated — the operational maturity available to run whatever you design.
State which properties you are prepared to give up. A design that claims to sacrifice nothing has simply not identified where it will fail.
Scaling: up, out, and along which axis
Scaling up — a bigger machine — is underrated. It preserves single-node consistency, keeps the mental model simple, and needs no new failure handling. Its limits are a hard ceiling and a coarse failure domain, but a very large fraction of systems never approach that ceiling.
Scaling out — more machines — is unbounded in principle but introduces the network into places that used to be function calls, which means partial failure, ordering questions, and coordination cost. The critical precondition is statelessness in the horizontally scaled tier: if instances hold state that only they know about, adding one does not add capacity, it adds an inconsistency.
Splitting a system can be done along three independent axes, and confusing them is a common source of a design that does not help: by load (identical clones behind a balancer, which helps with volume), by function (different services doing different things, which helps with independent scaling and team ownership), and by data (partitioning the same function across subsets of the data, which is the only axis that helps when a single dataset is the bottleneck). A system that is bottlenecked on one dataset gains nothing from more clones.
Data: partitioning and replication
The dominant early decision is the data model, and it should be derived from access patterns rather than from an abstract entity diagram. Storage engines differ chiefly in which access patterns they make cheap; choosing one before knowing the queries is choosing blind.
Replication makes copies. Single-writer-with-followers is simple and gives read scaling plus failover, at the cost of replication lag — reads from a follower may not reflect a write you just made, which produces the classic bug where a user saves something and immediately does not see it. Multi-writer setups remove that bottleneck and introduce write conflicts you must resolve.
Partitioning splits the data itself. The whole game is the partition key: it must spread load evenly and keep data that is queried together on the same partition. Get it wrong and you get hot partitions — one shard saturated while the rest idle — or queries that must fan out to every partition and wait for the slowest, which makes tail latency worse as you add shards. Repartitioning a live system is one of the most expensive operations in this field, which is why the key deserves disproportionate design attention up front.
Denormalisation trades write complexity and storage for read speed, and is often right for read-dominated workloads. The cost is that the same fact now exists in several places and can diverge, so you need a mechanism — a change stream, a rebuild job, a reconciliation check — that makes divergence detectable rather than permanent.
Consistency and transactions
When a network partition occurs, a distributed system must choose between refusing to answer and answering with possibly-stale data. That is the substance of the CAP result, and the useful part is what it implies when there is no partition: a system still trades latency against consistency, because stronger guarantees require more coordination, and coordination costs round trips. Both halves of that tradeoff are design choices you make deliberately or inherit by accident.
In practice you are choosing a point on a spectrum: strong guarantees where every read sees the latest write and the system pays in coordination and availability; read-your-own-writes, which is frequently what users actually notice and is much cheaper than full strength; and eventual convergence, which is cheapest and requires the application to tolerate — and the interface to present — temporarily disagreeing views.
Different data in one system warrants different levels. Financial balances and inventory counts usually justify coordination. A view counter or a recommendation list usually does not. Applying the strongest level everywhere is a real cost paid for no benefit.
Across services, distributed transactions are usually the wrong tool; the durable alternative is a saga — a sequence of local transactions, each with a compensating action that undoes it. This forces you to answer, explicitly, what "undo" means for each step, and some steps (an email sent, a payment captured) cannot truly be undone, only compensated. That is a product conversation as much as a technical one.
Synchronous versus event-driven
Synchronous request/response is simple to reason about and debug: the caller knows the outcome. Its cost is temporal coupling — the caller's availability is bounded by the callee's, and a chain of synchronous dependencies multiplies failure probability and adds latency at every hop.
Asynchronous messaging removes temporal coupling: the producer commits a fact and moves on. Its costs are that the outcome is not immediately known, that ordering is only as strong as the broker guarantees and is usually per-partition rather than global, and that the flow of a single logical operation becomes distributed across components and therefore much harder to follow without tracing.
A practical dividing line: use synchronous calls when the caller genuinely cannot proceed without the answer, and asynchronous events when other parts of the system merely need to know something happened. Emitting a fact and letting interested parties react also avoids the pattern where every new consumer requires a change to the producer.
Whichever you choose, remember that the boundary between committing state and publishing an event is the single most common correctness hole in distributed systems — the dual-write problem, whose standard remedy is the transactional outbox described in backend development.
Reasoning about bottlenecks
Two durable results let you reason about capacity without guessing at numbers you do not have.
Little's law states that the average number of items in a stable system equals the arrival rate multiplied by the average time each spends there. It is unconditionally true for any stable system, which makes it the fastest sanity check available: concurrency, throughput, and latency are not three independent dials, and fixing any two determines the third. If you want lower latency at fixed throughput, you must reduce in-flight work.
Queueing behaviour under utilisation is the other. As a resource approaches full utilisation, waiting time does not rise proportionally — it rises sharply and without bound, because variability in arrivals means the queue never fully drains. This is why a system that is comfortable at moderate utilisation becomes unusable after a modest traffic increase, and why running any shared resource near its capacity ceiling is a design error rather than an efficiency win.
The practical consequence: find the single most constrained resource and reason about that, rather than optimising components that are not the constraint. Parallelising the parts of a workload that were never the bottleneck produces disappointment that is entirely predictable — the serial fraction sets the ceiling.
Designing for failure
At any meaningful scale, something is always failing. The design question is not how to prevent that but what fails when it does.
- Failure domains. Identify what a single failure takes out. Redundancy inside one domain — two instances that share a power supply, a rack, a zone, or a control plane — provides far less protection than the instance count suggests.
- Graceful degradation. Decide in advance which features may be shed to keep the core working. A search page that drops personalised ranking and still returns results is behaving correctly; one that returns an error because a recommendation service is down is not.
- Isolation. Bulkheads — separate pools, quotas, or instances per dependency or per tenant — stop one failing dependency or one heavy customer from consuming a shared resource entirely.
- Backpressure and load shedding. When arrivals exceed capacity, the system must reject or slow intake. Accepting everything into unbounded queues converts overload into memory exhaustion and then into a crash loop, and the restarted system faces a larger backlog than before.
- Idempotency everywhere. Retries are inevitable in a distributed system, so every operation that can be retried must be safe to apply more than once.
- Bounded retries with jitter. Synchronised retries after a common failure produce a thundering herd that keeps the recovering dependency down.
The category of outage worth naming is metastable failure: a system that remains in a degraded state after the original trigger is gone, because the work generated by the degradation — retries, cache misses, reconnections — is now itself enough to sustain the overload. Backpressure, bounded retries, and the ability to shed load are what break that loop.
Operability
A design that cannot be observed cannot be operated, and operability should be part of the design rather than added afterwards. That means correlation identifiers that follow a request across every component; latency measured as a distribution rather than an average, since the mean hides exactly the failures users notice; explicit service objectives that define what "working" means; and health signals that distinguish "this instance is wedged" from "this instance should not receive traffic right now".
Design for deployment too. Rolling deploys mean two versions of your code run simultaneously against the same data, so every schema and message-format change must be compatible with both — the expand-and-contract discipline is not an optimisation, it is what makes zero-downtime deployment possible at all.
Common failure modes
- The distributed monolith. Services that must be deployed together: all of the operational cost of distribution, none of the independence.
- Premature partitioning. Sharding a dataset that a single well-indexed node would have served for years, buying irreversible complexity to solve a problem you did not have.
- The wrong partition key. Hot shards, or queries that fan out to everything and inherit the slowest response.
- The dual write. State and event committed separately; they diverge and stay diverged.
- Unbounded growth. A table, queue, or log with no retention policy, which works until it abruptly does not.
- Coordination on the hot path. A lock, a leader, or a strongly-consistent read in the most frequent operation, silently capping throughput.
- Retry storms. Independent retry logic at several layers, multiplying load precisely when the system is weakest.
- A single shared control plane. Every service depending on one configuration, discovery, or authentication component that has no independent failure story.
When not to distribute
Distribution is a cost paid continuously: network failure handling, partial failure, version skew, distributed debugging, more infrastructure, and more people needed to keep it running. It is worth paying when you have a real constraint that cannot be met another way — a component that must scale independently, a failure domain that must be isolated, a data boundary that must be enforced, or teams that are genuinely blocked on each other's release cadence.
It is not worth paying for organisational tidiness, for resume-shaped reasons, or in anticipation of a scale you have no evidence you will reach. A well-structured single deployable with clean internal module boundaries handles more load than most teams expect, is dramatically easier to debug, and — because the boundaries are already drawn — can be split later when a specific constraint justifies it. Designing internal seams early is cheap and reversible. Distributing early is expensive and is not.
Related reading
Backend Development
What backend engineering is really about: request lifecycle, where state lives, concurrency and connection pools, transactions, background work, caching, failure handling, and the operational habits that keep a service correct under load.
API Development
How to design, evolve, secure, and operate an HTTP API: protocol styles, contract design, error and pagination models, idempotency, versioning, rate limiting, and the failure modes that break clients.
Mobile Engineering
What makes mobile different from web and server work: irreversible releases, old versions in the wild, offline and sync, constrained resources, permissions, and the release engineering that keeps a shipped app recoverable.
Message Queue Architecture
Design reliable message-driven systems using queues and event streaming. Covers queue selection, delivery guarantees, dead letter handling, consumer patterns, and production monitoring.
Caching Strategies: Serving Data at the Speed of Memory
Implement caching that reduces latency from hundreds of milliseconds to single-digit milliseconds. Covers cache placement, invalidation strategies, cache-aside vs write-through patterns, distributed caching with Redis, CDN caching, and the pitfalls that turn your cache into a source of stale data and subtle bugs.
Request Coalescing
Production engineering guide for request coalescing covering patterns, implementation strategies, and operational best practices.
Distributed Tracing: Following Requests Across Service Boundaries
Implement distributed tracing to debug latency, identify bottlenecks, and understand request flow across microservices. Covers OpenTelemetry, trace propagation, span design, sampling strategies, and integrating traces with logs and metrics for full observability.