The idea: both split a table by a key, only one splits the machine

A table that grows too big is split into pieces by a key. The words for it get mixed up, so the page uses them in the most common way:

  • Partitioning: the pieces (partitions) stay inside one database server. It is a storage layout: one process, one set of CPUs and disks, one transaction log. PostgreSQL's declarative partitioning, MySQL's PARTITION BY and Oracle partitioning are this.
  • Sharding: the pieces (shards) are on different servers, each with its own CPU, memory, disk and failures. Something must route every query to the right server: a proxy (Vitess vtgate, MongoDB mongos, the Citus coordinator) or a library inside the application.

Sharding is sometimes called horizontal partitioning across servers, and systems such as Kafka and Cassandra call their shards "partitions". The difference that matters is not the word but whether the pieces share a machine. (Vertical partitioning, splitting columns into separate tables, is a different thing again and not shown.)

The page runs the same table, orders(order_id, user_id, month, total) with 12 rows, both ways. On the left, db1 partitions it by month (PARTITION BY RANGE). On the right, three servers shard it by user: shard = hash(user_id) % 3, with hash(u) = u to keep the arithmetic visible. Every button runs on both sides at once, in lockstep, and a cost line under each side counts servers, partitions or shards visited, network messages and steps. Cells are coloured by user, so you can see users scattered over months on the left and months scattered over users on the right.

Left: one server db1 holds the 12 orders in three monthly partitions, Jul, Aug and Sep, chosen by the planner. Right: a router sends each order by user_id % 3 to one of three servers: S0 holds users 3 and 6, S1 users 1 and 4, S2 users 2 and 5.
Partitioning splits the table inside one server; sharding spreads it over several servers behind a router.

Partitioning inside one server

A partitioned table is a parent with no storage of its own and one child table per partition. PostgreSQL offers three schemes:

SchemeExampleTypical use
RANGEFOR VALUES FROM ('2026-09-01') TO ('2026-10-01')time series, logs, orders: one partition per day, week or month
LISTFOR VALUES IN ('EU', 'UK')regions, tenants, status values
HASHFOR VALUES WITH (MODULUS 4, REMAINDER 1)spreading rows evenly when there is no natural range

On INSERT the server routes each row to its partition in memory (tuple routing). A row that fits no partition is an error, no partition of relation "orders" found for row, unless a DEFAULT partition exists; Demo: insert: planner vs router hits it with an October row. On SELECT the planner compares the WHERE clause with the partition bounds and skips partitions that cannot match: partition pruning. A query on the partition key opens one partition; a query on another column opens them all.

What partitioning buys: smaller indexes and tables that fit in memory better, pruning for queries on the key, maintenance per partition (VACUUM, reindex), and above all cheap removal of old data: DETACH PARTITION + DROP TABLE deletes a month by unlinking files, without reading a row (Demo: drop old data). What it does not buy: more CPU, memory, disk or write throughput than one server has, or any protection when that server fails.

Sharding across servers

Each shard is an ordinary database holding a subset of the rows. The shard key decides where a row lives; the router maps a key to a shard, either by a formula (hash(key) % N, ranges of the key) or by a lookup table (a directory). A query that carries the shard key goes to one shard. A query that does not must be sent to every shard and the answers merged: scatter-gather. Aggregates are pushed down, so each shard sends one partial result (Demo: aggregate over everything): SUM and COUNT combine directly, AVG is sent down as SUM and COUNT, a median cannot be combined this way.

What sharding buys: capacity that grows with the number of servers, for storage, memory and above all writes, and failures limited to one shard's rows. What it costs: a router, network hops on every query, fan-out for queries on other columns, transactions and joins across shards, and moving data when servers are added.

The key decides which queries are cheap

QueryLeft: partitioned by month, one serverRight: sharded by user, three servers
WHERE user_id = 3no pruning: every partition's index is probed, on one serverone shard
WHERE month = 'Sep'pruning: one partitionevery shard (scatter-gather)
SUM(total)one server reads everythingevery shard reads its part, in parallel
a transaction on two usersalways locallocal if both users share a shard, else two-phase commit
delete a monthdrop a partition, instantrow-by-row delete on every shard

Choose the shard key from the queries and transactions that must be fast, usually the tenant, customer or user, so that the data one request needs lives together. Demo: query with and without the key shows both sides winning one query and losing the other.

Grid of two queries against two layouts. WHERE month = Sep: the partitioned table visits only the Sep partition, the sharded table must ask all three shards. WHERE user_id = 3: the partitioned table probes every partition, the sharded table asks only S0.
Each layout makes the query on its own key cheap and the query on the other column expensive.

Transactions across shards: two-phase commit

On one server, any transaction is atomic through its single write-ahead log, however many partitions it touches. Across shards, each shard can only commit its own part, and if one commits while the other fails the data is wrong. Two-phase commit fixes that: the coordinator asks every participant to PREPARE (make the change durable but not visible, keep the locks, promise to commit), and only when all say yes does it send COMMIT PREPARED. It costs two extra round trips, locks are held across them, and if the coordinator dies between the phases the participants must wait for it with their locks held. Demo: transaction across two users runs one transaction inside S1 (users 1 and 4) and one across S1 and S2 (users 1 and 2). Many sharded systems avoid cross-shard transactions by choosing the key so they do not happen, or offer them only with weaker guarantees.

Adding a server: why not hash % N

With shard = hash(key) % N, changing N changes the answer for most keys: going from 3 to 4 shards on the page moves 4 of 6 users, and with many shards nearly every key moves (Demo: add a server with hash % N). Real systems put a fixed, larger number of buckets between keys and servers: the key picks a bucket (which never changes), and a small table maps buckets to servers. Adding a server hands it some buckets; only their rows move (Demo: add a server with a bucket map, 2 of 6 users). Redis Cluster's 16384 hash slots, Vitess keyspace ranges, Citus shard groups, Cassandra's token ring and consistent hashing are all versions of this. (see Consistent Hashing for the ring, virtual nodes and fixed slots side by side). Moving is done online: copy the bucket, stream the changes made meanwhile, then switch the router.

Before: S0 holds users 3 and 6, S1 users 1 and 4, S2 users 2 and 5. After adding S3 with hash % 4, users 3, 4, 5 and 6 all change servers; with a bucket map, S3 takes over the buckets of users 4 and 5 and nobody else moves.
Re-hashing with a new N moves most keys; a fixed bucket map moves only the buckets handed to the new server.

Failures and hot spots

A partitioned table is exactly as available as its server: when db1 goes down, every partition goes with it. When a shard goes down, only its keys are unavailable; queries for other users go on, but any query that needs every shard (a month, a sum) cannot give a complete answer (Demo: a server fails). Either way, each server still needs replicas and failover for its own data; see the Redis Sentinel and Redis Cluster pages.

Sharding spreads load only as well as the key spreads it. A single very active user sends all of its writes to one shard while the others idle (Demo: hot user, watch the request counters). Sharding by time would be worse: every new row would land on the newest shard. Partitioning by time is fine in comparison, because all partitions share one server anyway.

Summary

PartitioningSharding
Where the pieces liveone servermany servers
Who routesthe planner, in memorya router or client library, over the network
Capacityone machine (scale up)grows with servers (scale out)
Transactions and joinsnormallocal per shard; across shards 2PC or not at all
Failureall or nothingone shard's data
Adding a piecemetadatadata movement
Best keyoften time, for pruning and retentionthe entity requests revolve around (tenant, user)
Operational costlowhigh: routing, resharding, schema changes on N databases, backups of N databases

They are not alternatives but layers. A common design shards by tenant or user for capacity, and partitions each shard's large tables by time for pruning and retention. And a well-sized single server with partitioning goes a long way: shard when one machine really is not enough.

What the page leaves out

Replicas of each shard and failover, global secondary indexes, cross-shard joins, reference tables copied to every shard, directory-based sharding services, online resharding details (copy, catch-up, cut-over), consistent hashing rings and virtual nodes, sub-partitioning, partition-wise joins and aggregates, and the real hash functions.