Skip to main content

Sharding

Updated 2 min read

Share this page

Send the link, quote the definition with a link back, or show it as a card on your own site.

https://softwaredictionary.org/terms/sharding

In short

Sharding is a way of scaling a database by splitting its data across several servers, called shards, so each one stores and handles only part of the total.

What is database sharding?

Sharding breaks one large database into smaller pieces, called shards, that live on separate servers. Each shard holds a subset of the rows, for example users A to M on one server and N to Z on another, and all shards share the same schema. Together they behave like one logical database that can store more data and handle more traffic than any single machine.

A shard key decides where each row goes. With range-based sharding, rows are split by ranges of the key, such as dates or ID ranges; with hash-based sharding, the key is run through a hash function so rows spread evenly; and with directory-based sharding, a lookup table maps each key to its shard. The application, a routing layer, or the database itself uses the shard key to send each query to the right server.

Think of a large library that splits its collection across several buildings by the author's last name: each building is smaller and less crowded, but you have to know which one to visit. Sharding is used by large web applications and is built into or available for systems such as MongoDB, Apache Cassandra, MySQL with Vitess, and PostgreSQL with Citus.

Sharding is often confused with replication. Replication copies the same data to several servers for availability and read scaling, while sharding splits different data across servers to scale writes and storage, and large systems usually combine both by replicating each shard. Because queries and transactions that span shards are slower and harder, and a poor shard key can create hot spots where one shard gets most of the traffic, teams usually shard only after indexing, caching, and bigger hardware are no longer enough.

At a glance

Sharding: the users table is split by user_id across three database servers, with shard 1 holding ids 1 to 999, shard 2 ids 1000 to 1999 and shard 3 ids 2000 to 2999; a router reads the shard key of a query for user_id 2481 and sends it only to shard 3.one table, split into shardsQueryuser_id = 2481Routershard key: user_idShard 1user_id 1 – 999Shard 2user_id 1000 – 1999Shard 3user_id 2000 – 2999
Each server holds only part of the rows. The shard key tells the router which one to ask, so the other shards never see the query.

Key takeaways

  • Sharding splits a database's rows across multiple servers.
  • A shard key determines which shard stores each row.
  • Common strategies are range-based, hash-based, and directory-based sharding.
  • Sharding scales writes and storage; replication copies data for availability and reads.
  • Cross-shard queries and a poorly chosen shard key are the main pitfalls.

Example

Routing queries with hash-based shardingjavascript
// Pick a shard from the user ID, so all of a user's rows live together
import { createHash } from "node:crypto";

const shards = [dbShard0, dbShard1, dbShard2, dbShard3];

function shardFor(userId) {
  const hash = createHash("md5").update(String(userId)).digest();
  return shards[hash.readUInt32BE(0) % shards.length];
}

const db = shardFor(42);
const orders = await db.query("SELECT * FROM orders WHERE user_id = $1", [42]);

Readers ask

What is the difference between sharding and replication?

Sharding splits data so each server holds a different part of it, which scales storage and writes. Replication copies the same data to several servers, which improves availability and read capacity. Many production systems use both.

What is the difference between sharding and partitioning?

Partitioning is the general idea of splitting a table into parts, often within a single database server. Sharding is horizontal partitioning across multiple servers, so each part runs on its own machine.

When should you shard a database?

Usually only when a single server can no longer handle the data size or write load after you have tried indexing, caching, query tuning, read replicas, and larger hardware. Sharding adds lasting complexity, so it is rarely the first scaling step.

Often compared

See also

Spotted a mistake or something missing on this page?Suggest an edit

Read a random page
Open today's review
Switch to the dark theme
Read this page in Türkçe

More

Settings