As applications grow, databases must handle increasing amounts of data and more concurrent requests. A database that performs perfectly with a few hundred users might struggle when thousands or millions of users start interacting with it.

This raises an important question: How do we scale a database when a single server is no longer enough?

Two common approaches are vertical scaling and horizontal scaling. While both aim to improve database capacity and performance, they work differently and introduce different challenges.

Understanding these approaches also helps explain why traditional relational databases, such as PostgreSQL and MySQL, are often considered harder to scale horizontally than databases like MongoDB.

What Is Database Scaling?

Database scaling is the process of increasing a database system's capacity to handle growing workloads, including larger datasets, more queries, and higher numbers of concurrent users.

There are two primary approaches: vertical scaling (scaling up) and horizontal scaling (scaling out).

Vertical Scaling

Vertical scaling means increasing the resources of an existing database server, such as adding more CPU, RAM, or faster storage.

For example, imagine a PostgreSQL database running on a server with 4 CPU cores and 8 GB of RAM. As the application grows, we upgrade the server to 16 CPU cores and 64 GB of RAM.

The database continues running on a single server, but that server now has more resources to process requests.

Advantages: Vertical scaling is relatively straightforward because it usually requires fewer architectural changes. Applications can continue using the same database connection and existing queries.

Limitations: Hardware resources are finite. Eventually, upgrading a server becomes expensive or reaches the maximum capacity supported by the infrastructure.

Horizontal Scaling

Horizontal scaling means adding more servers to distribute database workloads instead of relying on a single powerful machine.

For example, rather than having one database server handle 100,000 requests, we introduce additional database servers to share the workload.

However, adding database servers does not automatically distribute requests or data. The database architecture must determine how those servers work together.

Two common techniques used in horizontal database scaling are replication and sharding.

Database Replication: Multiple Copies of the Same Data

Replication involves maintaining copies of the same database across multiple servers.

Imagine an e-commerce application with two tables: Users and Orders.

In a typical primary-replica architecture, both servers maintain copies of these tables.

Server A (Primary):

  • Users table
  • Orders table
  • Handles write operations

Server B (Replica):

  • Users table
  • Orders table
  • Handles read operations

When a customer creates an order, the application writes the new record to the primary database. The primary's changes are then replicated to the secondary server.

Read requests can be distributed across replicas, reducing the workload on the primary.

Why Replication Is Useful

Read Scalability: Applications with many read requests can distribute queries across multiple replicas. For example, displaying product catalogs, retrieving user profiles, and browsing order histories can use read replicas when their consistency requirements allow it.

High Availability: If the primary database fails, a suitable replica can be promoted to become the new primary, depending on the failover configuration.

Workload Isolation: Reporting and analytics queries can be directed to replicas to reduce their impact on the primary database.

The Limitation of Replication

Replication does not necessarily solve write scalability.

In a traditional single-primary configuration, all insert, update, and delete operations still go through one primary database.

Even if we add ten read replicas, the primary may remain a bottleneck when the application experiences heavy write traffic.

Additionally, replication can be asynchronous, meaning replicas may temporarily contain older data than the primary.

This brings us to another approach: sharding.

Database Sharding: Distributing Data Across Servers

Unlike replication, where servers maintain copies of the same data, sharding distributes different portions of a dataset across multiple database servers.

Each portion is called a shard.

Consider an e-commerce application with millions of users and orders. Instead of keeping all records on one server, we distribute them based on user IDs.

For example:

Shard A:

  • Users with IDs 1–1000
  • Orders belonging to those users

Shard B:

  • Users with IDs 1001–2000
  • Orders belonging to those users

Shard C:

  • Users with IDs 2001–3000
  • Orders belonging to those users

These ranges are simplified examples. Production databases may use hash-based distribution, range-based distribution, or other partitioning strategies.

When the application requests orders for user 1500, the system identifies the appropriate shard and retrieves the records from Shard B.

Because the data is distributed, different servers can process different requests simultaneously.

How Does the Database Know Which Shard to Query?

A sharded database needs a mechanism to determine where data belongs.

This typically involves a shard key, which is a field used to determine how records are distributed.

For example, using user_id as a shard key allows records to be distributed according to their associated users.

A query router or application-level routing mechanism identifies the appropriate server based on that key.

Choosing a good shard key matters because a poor distribution strategy can cause uneven workloads, where one shard handles most requests while others remain underutilized.

Benefits of Sharding

Write Scalability: Different shards can independently process writes for the data they own, increasing overall write capacity when workloads are well distributed.

Storage Scalability: Large datasets can be distributed across multiple machines rather than requiring one server to hold everything.

Workload Distribution: Queries that target specific shards can be processed independently, reducing contention between unrelated operations.

However, these benefits come with additional architectural complexity.

Why Are Traditional Relational Databases Harder to Shard?

Relational databases are built around structured relationships, joins, constraints, and transactions.

These features are valuable because they help maintain data integrity. However, distributing relational data across independent servers introduces additional challenges.

Cross-Server JOIN Operations

Consider two relational tables: Users and Orders.

A query might retrieve a customer's orders using a JOIN:

SQL
SELECT users.name, orders.amount
FROM users
JOIN orders ON users.id = orders.user_id
WHERE users.id = 1500;

When both tables and the relevant records are on the same server, the database can execute the JOIN locally.

But what happens if the required records are distributed across different servers?

The system may need to fetch data over the network, coordinate query execution, and combine results.

This introduces network overhead and additional processing costs.

A common solution is to organize related data so that records frequently accessed together are stored on the same shard. However, not every query can benefit from the same distribution strategy.

Distributed Transactions

Relational databases support transactions that allow multiple database operations to succeed or fail as a unit.

Imagine a customer purchasing a product. The application must create an order and reduce the available inventory.

When both operations happen inside one database server, a local transaction can coordinate them efficiently.

But if the order is stored on Shard A and the inventory record is stored on Shard B, the transaction involves multiple independent servers.

The database must coordinate the operations and determine whether the transaction can safely commit.

Protocols such as two-phase commit can help coordinate distributed transactions, but they introduce communication overhead and more complicated failure handling.

Data Consistency and Relationships

Traditional relational databases commonly enforce constraints such as foreign keys to prevent invalid relationships.

For example, an order should not reference a user who does not exist.

Enforcing such relationships becomes more difficult when the related records reside on different servers.

A distributed database may need additional coordination, restrict certain cross-shard constraints, or move some integrity checks into the application layer.

This does not mean relational databases cannot scale horizontally. It means preserving relational guarantees across multiple machines requires additional engineering.

Does SQL Support Sharding?

Yes. Sharding is not exclusive to MongoDB or NoSQL databases.

Several SQL databases and distributed SQL systems support sharding, although their implementations differ.

PostgreSQL: PostgreSQL supports native table partitioning, but partitioning alone does not automatically distribute data across independent servers. Distributed PostgreSQL solutions, such as Citus, can provide sharding across multiple nodes.

MySQL: MySQL applications can implement sharding through application-level routing or solutions such as Vitess.

CockroachDB: CockroachDB is a distributed SQL database designed to distribute data and workloads across multiple nodes while supporting SQL transactions.

YugabyteDB: YugabyteDB provides a distributed SQL architecture with automatic data distribution across database nodes.

The important distinction is that traditional PostgreSQL and MySQL deployments do not automatically distribute writes across independent servers merely because more servers have been added.

Why Is MongoDB Commonly Associated with Sharding?

MongoDB was designed with built-in support for distributing document data across shards.

Its sharded cluster architecture includes several important components.

Shards: Store different subsets of the distributed dataset. In MongoDB's standard sharded architecture, each shard is deployed as a replica set.

Query Routers: MongoDB uses mongos processes to route application queries to the appropriate shards.

Configuration Servers: Maintain metadata about the sharded cluster and where data is distributed.

Shard Keys: Determine how documents are distributed across the cluster.

MongoDB also includes mechanisms for balancing data across shards, helping distribute storage and workloads as the dataset grows.

However, MongoDB is not immune to distributed database challenges.

Queries involving multiple shards, distributed transactions, and poorly chosen shard keys can still create performance bottlenecks.

MongoDB makes sharding an integrated database capability, but that does not mean every MongoDB application needs sharding or that sharding is always easy.

Replication vs. Sharding: Understanding the Difference

Although replication and sharding both use multiple database servers, they solve different problems.

Feature Replication Sharding
Data storage Copies of the same data Different portions of data
Primary purpose Read scaling and availability Data and workload distribution
Read scalability Yes, through read replicas Yes, for distributable queries
Write scalability Limited with a single primary Can distribute writes
Complexity Usually lower Usually higher
Typical challenge Replication lag and failover Data distribution and cross-shard operations

An important detail is that replication and sharding can be combined.

For example, an application might distribute its data across three shards, with each shard having its own primary and replica servers.

This architecture allows the system to distribute writes across shards while using replication to improve availability and potentially serve additional read requests.

When Should You Consider Sharding?

Sharding is a powerful technique, but it should not be the first solution whenever a database becomes slow.

Many database performance problems can be addressed without introducing distributed storage.

Before considering sharding, it is often worth examining query performance, indexing strategies, connection pooling, caching, table partitioning, and server resources.

Read replicas may also be sufficient for applications with heavy read traffic but relatively moderate write workloads.

Sharding becomes more relevant when a single database server cannot reasonably accommodate the required data volume or write throughput, even after appropriate optimization.

The decision should depend on actual workload requirements rather than assumptions about SQL or NoSQL performance.

Conclusion

Database scaling is not simply about adding more powerful servers or introducing additional database instances. It involves choosing an architecture that matches an application's workload, consistency requirements, and expected growth.

Vertical scaling increases the capacity of an existing server, while horizontal scaling distributes workloads across multiple servers.

Replication helps improve read scalability and availability by maintaining copies of data. Sharding distributes different portions of the dataset, allowing storage and write workloads to scale across multiple machines.

Traditional relational databases can support horizontal scaling, but distributing relational operations introduces challenges involving JOINs, transactions, constraints, and consistency.

MongoDB provides integrated sharding capabilities, while PostgreSQL, MySQL, and distributed SQL databases offer alternative approaches to achieving similar goals.

Ultimately, the challenge of horizontal scaling is not SQL versus NoSQL. It is managing data distribution, coordination, and consistency across independent servers.

Understanding these trade-offs helps developers make better architectural decisions and, just as importantly, recognize when a simpler single-database deployment is already sufficient.

References