Shared-Nothing Architecture: The Foundation of Linear Scalability

Philip Rehberger Aug 18, 2026 7 min read

Eliminate shared state to scale horizontally without contention. Covers partitioning, replication, and routing.

Shared-nothing architecture is the principle behind almost every system that scales horizontally. The idea is older than the buzzword: each node in the system has its own CPU, memory, and storage, and shares none of them with any other node. Nodes coordinate through messages, not through shared state. The result is a system whose throughput grows linearly as you add nodes — at least until something else becomes the bottleneck.

The term comes from Michael Stonebraker's 1986 paper, which contrasted shared-nothing with shared-memory and shared-disk database architectures. Forty years later, "just add nodes" is the default scaling strategy for the entire web, and almost every system that does it is shared-nothing under the hood.

The Three Architectural Patterns

Shared-memory. Multiple processors access a common pool of memory. Scales until the memory bus is saturated, which happens quickly.

Shared-disk. Multiple nodes have separate memory but share storage. Common in older database clusters (Oracle RAC, IBM DB2 PureScale). The storage becomes a coordination point and limits how far you can scale.

Shared-nothing. Each node owns its memory and its storage. Scaling is bounded by the network and by whatever coordination is required between nodes.

The first two architectures get you to maybe a few dozen nodes before something gives. Shared-nothing systems scale to thousands of nodes because there is no central resource to contend for.

What "Shared Nothing" Actually Requires

The label is easy to claim and harder to earn. A system is shared-nothing only if:

  • Each node has its own local storage. No SAN, no NAS, no cluster filesystem.
  • Each node owns a slice of the data. Sharding or partitioning is the rule, not the exception.
  • Nodes coordinate by message passing, not by reading or writing the same memory or disk.
  • Adding a node requires only network connectivity — no special hardware, no manual data placement.

If two of your "nodes" share a network filesystem, you are not shared-nothing. If your nodes can read each other's caches, you are not shared-nothing. These are coordination points and they will eventually limit you.

How Sharding Implements Shared-Nothing

The mechanical answer to "how do we share nothing" is sharding: each node owns a partition of the data, and routing decides which node handles which request.

Request for user 12345
      │
      ▼
+------------------+
| Router / Proxy   |   shard_key = hash(user_id) % N
+--------+---------+
         │
   ┌─────┼─────┐
   ▼     ▼     ▼
+----+ +----+ +----+
| N0 | | N1 | | N2 |
+----+ +----+ +----+
 each owns its own slice; no node reads another's data

The router computes the shard from the request and forwards it. Each node serves its slice. There is no cross-node data access on the hot path.

Sharding strategies range from simple hashing to consistent hashing (which minimizes rebalancing when nodes are added or removed) to range-based partitioning (good for ordered queries, bad for hotspots).

The hard part is not the sharding mechanic. The hard part is choosing the shard key. A bad key produces hotspots (one shard takes all the traffic) or cross-shard queries (defeating the point of shared-nothing).

Cross-Shard Queries: Where the Pattern Gets Tested

A pure shared-nothing system is great for queries that hit one shard. The moment you need data from multiple shards, the architecture has to make a choice.

Scatter-gather. Send the query to all shards in parallel, collect the results, combine them. Works for small numbers of shards; latency degrades as shards multiply.

Cross-shard joins via duplication. Denormalize so each shard has the data it needs. Eats storage, gives back query simplicity.

A second store optimized for cross-shard queries. Build a search index or a denormalized read model alongside the shared-nothing source of truth.

The "right" answer depends on how often the cross-shard query happens and how stale the result can be. For real-time exact answers, scatter-gather is the only choice. For approximate or slightly-stale answers, a separate read model is usually better.

Shared-Nothing in Common Systems

The pattern shows up everywhere once you know to look for it:

System How it is shared-nothing
Cassandra Each node owns a token range; replicas are explicit and bounded
DynamoDB Each partition lives on a primary and N replicas; no shared storage
Elasticsearch Each shard lives on one node; cross-shard search is scatter-gather
Stateless web tier Each web node holds no state; session data lives in a cache or DB
Kafka Each partition has one leader; producers and consumers know where to go

The web tier example is instructive: most engineers do not call this "shared-nothing," but that is exactly what stateless application servers are. They scale linearly because they share nothing — every request can be handled by any node, and adding a node just adds capacity.

What the Pattern Costs

Shared-nothing buys scalability at the cost of complexity:

  • Routing is essential. Every request has to find the right node. This adds a hop and a potential failure mode.
  • Resharding is hard. Adding nodes means moving data. Done well, this is a controlled rebalance. Done badly, it is downtime.
  • Cross-shard transactions are expensive. Whatever consistency you needed within a shard either disappears or becomes a distributed transaction.
  • Hotspots are real. If your shard key has a heavy hitter (a popular user, a high-traffic product), one node carries the load while others sit idle.

These costs are why most teams should start with a single, larger node — vertical scaling — and only adopt shared-nothing when they have outgrown what one node can do.

The Database Caveat

The classic vertical-then-horizontal scaling path on relational databases is: bigger server, then read replicas, then sharding. Sharding is where you go shared-nothing on the write path, and it is one of the highest-effort engineering projects most companies undertake. Frameworks like Vitess and Citus do a lot of the heavy lifting, but the data-modeling work — picking shard keys, redesigning cross-shard queries — is yours.

The advice that has aged well: do not shard until you have to. A modern relational database on a single large instance handles workloads that would have required shared-nothing fifteen years ago. Test the limits of vertical scaling and read replicas first.

Choosing the Shard Key

When you do shard, the shard key is the single most important decision. Good shard keys:

  • Distribute load evenly across nodes
  • Match the access pattern (most queries hit a single shard)
  • Are stable (the key for a record does not change over its lifetime)
  • Are present in nearly every query the application makes

For most multi-tenant SaaS products, tenant ID is a strong shard key. Each tenant's data is naturally isolated, queries almost always include the tenant ID, and loads are balanced once you have enough tenants. For consumer products, user ID often plays the same role.

Why It Matters

Shared-nothing is the architecture that lets web-scale systems exist. Twitter, Netflix, every cloud database, every CDN — under the hood, the pieces that scale are shared-nothing. The pattern is not glamorous, and the details are messy, but it is the foundation that supports almost everything else.

For most teams, the lesson is to recognize when shared-nothing patterns are needed and apply them deliberately when the time comes — not to copy the architecture of systems several orders of magnitude larger than the one you are actually building.


Sizing a system whose load is starting to outgrow vertical scaling? We help teams figure out which boundary to shard on, when, and what to leave alone. scopeforged.com

Share this article

Related Articles

Need help with your project?

Let's discuss how we can help you build reliable software.