Short answers to what readers ask most about this topic.
01What is the difference between database partitioning and sharding?
Partitioning splits a table into smaller pieces inside one database server, and the database engine routes rows to the right piece. Sharding splits data across several servers, and your application or a routing layer decides which server owns a row. Partitioning improves scan size and retention, while sharding adds capacity.
02Is sharding the same as horizontal partitioning?
Sharding is a form of horizontal partitioning, meaning different rows go to different places. The difference is that horizontal partitioning usually stays inside one database instance, while sharding applies the same split across multiple instances. Vertical partitioning is something else: it splits columns, not rows.
03When should I partition a PostgreSQL table?
Partition when a table is very large and most queries filter on one key, typically a date. The Postgres documentation gives a rule of thumb that the table should be bigger than the server's physical memory. The biggest practical win is cheap retention: detaching or dropping an old partition avoids a slow bulk DELETE and its VACUUM cost.
04Does PostgreSQL partitioning make queries faster automatically?
Only queries that filter on the partition key benefit, because the planner can prune partitions it proves cannot match. A query that filters on another column probes every partition, as an EXPLAIN on a table partitioned by date and filtered by branch shows. Too many partitions can also lengthen planning time.
05When is sharding unavoidable?
Sharding becomes unavoidable when write volume or data size exceeds what the largest single server you can run can handle, even after indexes, read replicas, caching and partitioning. It can also be required for data residency, where a region's data must stay on servers in that region. Expect cross-shard queries, hot shards and resharding work in return.
Partitioning splits one large table into smaller pieces inside a single database server, so queries scan less and old data can be dropped cheaply. Sharding spreads data across several servers, each owning a slice chosen by a shard key. Partition first for retention and scan size; shard only when one machine cannot hold the writes or the data.
The two words get used as synonyms in design reviews, and the confusion is expensive, because the two techniques fix different problems. A partitioned Postgres table is still one database on one server. A sharded system is several databases, and your application suddenly has to know which one holds a row.
This post is the concept comparison. For the hands-on Citus and rollout side of the topic, see the existing post PostgreSQL Sharding Strategy; here I stay on the difference, with Postgres DDL I ran myself, a TypeScript shard router, and a checklist for what to try before either. The facts about Postgres come from its official partitioning documentation.
What is the difference between partitioning and sharding?
Both split a dataset by rows. The difference is where the pieces live. Wikipedia describes a shard as a horizontal partition of data in a database, and notes that horizontal partitioning usually happens within a single database instance while sharding applies the same row split across multiple instances.
Question
Partitioning (Postgres declarative)
Sharding
Where do the pieces live?
As child tables inside one database on one server
In separate databases, normally on separate servers
Who decides which piece holds a row?
The database engine, from the partition bounds
Your routing layer or application, from the shard key
What do transactions and joins look like?
Normal single-node transactions and joins across partitions
Cross-shard transactions and joins are hard or unavailable
What problem does it solve?
Scan size, index size, cheap retention, bulk load and delete
A write load or data size one server cannot carry
Does it add hardware capacity?
No, it is still one machine
Yes, each shard brings its own CPU, memory and disk
What does it cost to change later?
Attach or detach a partition, no application change
Resharding moves data and usually touches the application
The rest of the post follows the rows of that table. If a row does not describe your problem, you probably do not need the technique on that side of it.
How does table partitioning work in PostgreSQL?
Postgres declarative partitioning offers three methods: range partitioning by non-overlapping ranges of a key, list partitioning by explicit key values, and hash partitioning by a modulus and remainder. Time-series data such as orders or logs almost always wants range. The following DDL is modelled on the order log of a carwash POS, and I ran it on a local PostgreSQL 16.15.
-- Range partitioning by month on a table shaped like a carwash POS order log.
-- The primary key MUST include the partition key (created_at), or Postgres refuses it.
CREATE TABLE wash_orders (
id bigint GENERATED ALWAYS AS IDENTITY,
branch_id int NOT NULL,
created_at timestamptz NOT NULL,
total_idr integer NOT NULL,
PRIMARY KEY (id, created_at)
) PARTITION BY RANGE (created_at);
-- The upper bound is exclusive, so adjacent months never overlap.
CREATE TABLE wash_orders_2026_09 PARTITION OF wash_orders
FOR VALUES FROM ('2026-09-01') TO ('2026-10-01');
CREATE TABLE wash_orders_2026_10 PARTITION OF wash_orders
FOR VALUES FROM ('2026-10-01') TO ('2026-11-01');
-- Created on the parent, it becomes one index per partition automatically.
CREATE INDEX ON wash_orders (branch_id, created_at);
-- Synthetic data: 40,001 rows, one every 2 minutes starting 2026-09-01.
INSERT INTO wash_orders (branch_id, created_at, total_idr)
SELECT 1 + (g % 5),
timestamptz '2026-09-01' + (g * interval '2 minutes'),
50000 + (g % 7) * 10000
FROM generate_series(0, 40000) AS g;
ANALYZE wash_orders;
Partition pruning is the payoff. With pruning enabled, which is the default through the enable_partition_pruning setting, the planner proves that a partition cannot contain rows matching the WHERE clause and leaves it out of the plan. Below is the real output from my run. The first query filters on the partition key and touches only the October partition. The second filters on branch_id alone, so nothing can be pruned and both partitions are probed.
-- Filter on the partition key: only October is scanned.
EXPLAIN (COSTS OFF)
SELECT count(*) FROM wash_orders
WHERE created_at >= '2026-10-01' AND created_at < '2026-10-08';
Aggregate
-> Seq Scan on wash_orders_2026_10 wash_orders
Filter: ((created_at >= '2026-10-01 00:00:00+07'::timestamp with time zone) AND (created_at < '2026-10-08 00:00:00+07'::timestamp with time zone))
-- Filter on something else: nothing can be pruned, both partitions are probed.
EXPLAIN (COSTS OFF)
SELECT count(*) FROM wash_orders WHERE branch_id = 3;
Aggregate
-> Append
-> Bitmap Heap Scan on wash_orders_2026_09 wash_orders_1
Recheck Cond: (branch_id = 3)
-> Bitmap Index Scan on wash_orders_2026_09_branch_id_created_at_idx
Index Cond: (branch_id = 3)
-> Bitmap Heap Scan on wash_orders_2026_10 wash_orders_2
Recheck Cond: (branch_id = 3)
-> Bitmap Index Scan on wash_orders_2026_10_branch_id_created_at_idx
Index Cond: (branch_id = 3)
Two honest caveats about this run: the data is synthetic and tiny, so I am showing plan shapes, not timings, and the planner chose a sequential scan for October only because 18,401 rows is too few for an index to win. The +07 in the plan is my machine's time zone. The lesson that carries over is the second query: a partitioned table only helps queries that filter on the partition key.
Why is dropping a partition cheaper than deleting rows?
Retention is the strongest reason to partition. The Postgres documentation says dropping a partition with DROP TABLE or ALTER TABLE DETACH PARTITION is far faster than a bulk operation, and that it avoids the VACUUM overhead a bulk DELETE leaves behind. For a table that keeps, say, 13 months of orders, the monthly cleanup becomes a metadata change instead of a long-running delete.
-- Wrong: a bulk DELETE touches every row, bloats the table and leaves VACUUM work behind.
DELETE FROM wash_orders WHERE created_at < '2026-10-01';
-- Right: detach the old month (CONCURRENTLY avoids the heavy lock on the parent),
-- archive it if you must, then drop it. No per-row work happens at all.
ALTER TABLE wash_orders DETACH PARTITION wash_orders_2026_09 CONCURRENTLY;
-- pg_dump -t wash_orders_2026_09 ... (optional archive step)
DROP TABLE wash_orders_2026_09;
The same documentation notes that a plain DROP TABLE needs an ACCESS EXCLUSIVE lock on the parent, and offers DETACH PARTITION, with a CONCURRENTLY form, as the option that is often preferable because it keeps the data available as an ordinary table for a backup first. The docs also state a rule of thumb for when this pays off: the table should be larger than the physical memory of the database server.
The primary key trap. A unique or primary key on a partitioned table must include every partition key column, so a key on id alone is rejected and you end up with primary key (id, created_at). That means the database cannot enforce uniqueness of id across months. If you need that guarantee, enforce it another way, for example with an identity column or a UUID that your application never reuses.
Unique and primary key constraints must include all partition key columns, and cannot use expressions or function calls in the partition key.
There is no exclusion constraint across the whole partitioned table, only on each leaf partition.
Planning stays fast up to a few thousand partitions, but only if typical queries let the planner prune all but a few of them. Daily partitions for ten years is about 3,650 partitions, which is already at the edge of that guidance.
Inserting a row that fits no partition raises an error, so something must create next month's partition before the month starts.
Why do horizontal and vertical partitioning confuse people?
Three overlapping vocabularies are in circulation, and search results mix them freely. Keep these definitions straight and most of the argument disappears.
Horizontal partitioning puts different rows in different tables. Both Postgres range partitioning and sharding are horizontal; they differ in the server, not in the shape of the split.
Vertical partitioning puts different columns in different tables, for example moving a rarely read description or blob column out of a hot table. It is also called row splitting, and it has nothing to do with sharding by itself.
Some systems use the word partition for what others call a shard. Wikipedia notes that MongoDB, Elasticsearch and SolrCloud use the word shard for what its article calls partitions, while Postgres uses partition for a child table on one server. Ask which server owns the piece before assuming anything.
A useful test when you read an architecture diagram: can one machine failing take only part of the table away? If yes, it is sharding. If the whole table goes down together, it is partitioning, whatever the box is labelled.
How does sharding work, and what does it cost?
Sharding needs a shard key, a column whose value decides the owning shard, and a routing layer that applies it. The sketch below assumes each shard holds the same schema with a tenant_id column, as in a multi-tenant ERP where one customer company is one tenant. It hashes the tenant id with sha256 and takes the remainder.
import { createHash } from "node:crypto";
import { Pool } from "pg";
// One Pool per shard. In production these point at different servers.
const SHARDS: Pool[] = [
new Pool({ connectionString: process.env.SHARD_0_URL }),
new Pool({ connectionString: process.env.SHARD_1_URL }),
new Pool({ connectionString: process.env.SHARD_2_URL }),
new Pool({ connectionString: process.env.SHARD_3_URL }),
];
// Do not use Math.random() or JS string hashCode: the result must be stable
// across processes and deploys, or a tenant's rows are written to one shard
// and looked up on another.
export function shardFor(tenantId: string, shardCount = SHARDS.length): number {
const digest = createHash("sha256").update(tenantId).digest();
return digest.readUInt32BE(0) % shardCount;
}
export function poolFor(tenantId: string): Pool {
return SHARDS[shardFor(tenantId)];
}
// Single-tenant query: one hop, one shard.
await poolFor(tenantId).query(
"SELECT * FROM wash_orders WHERE tenant_id = $1 AND created_at >= $2",
[tenantId, from],
);
// Cross-shard query: the application, not the database, must fan out and merge.
const parts = await Promise.all(
SHARDS.map((p) => p.query("SELECT count(*) AS n FROM wash_orders")),
);
const total = parts.reduce((sum, r) => sum + Number(r.rows[0].n), 0);
The router is the easy half. The hard half is changing the number of shards. With hash modulo N, adding a fifth shard to four changes the answer for most keys. The arithmetic and a check against 100,000 generated tenant ids are below.
// Growing from 4 shards to 5 with hash % N. A key keeps its shard only when
// h % 4 == h % 5. Over any 20 consecutive hash values that is h = 0,1,2,3:
// stay = 4 / 20 = 20% move = 16 / 20 = 80%
//
// Checked on 100,000 ids "tenant-0" .. "tenant-99999" with the sha256 router above:
// kept their shard: 19,962 (19.96%) moved: 80,038 (80.04%)
// A consistent-hash ring or a lookup table moves only about 1/N of the keys.
That is why a plain modulo router is a starting point, not a design. The post Consistent Hashing Explained covers the ring that limits movement to roughly one in N keys, and a lookup table that maps tenants to shards gives you the same freedom with the cost of one more thing to keep consistent. The Citus documentation describes the same idea from the database side, with a distribution column and co-located tables so that related rows land on the same node.
Choose the shard key so the most common query carries it. In a multi-tenant system that is tenant_id: a request for one tenant goes to one shard, joins between that tenant's tables stay local, and small shared tables such as price lists are copied to every shard, which Citus calls reference tables.
Hot shards: if one tenant or one key is far larger than the rest, its shard saturates while the others idle, and hashing cannot fix a single oversized key.
Cross-shard queries: a report across all tenants must fan out to every shard and merge the results in application code, as in the Promise.all above.
Cross-shard transactions: there is no single commit across servers, so you either avoid them or take on two-phase commit or sagas.
Operations: schema changes, backups and failover now run per shard, and every shard needs its own replica to be safe.
What should you try before partitioning or sharding?
Most slow tables are not a partitioning problem. Work down this list in order, because each step is cheaper and more reversible than the next.
Symptom
Try first
Why it comes before partitioning or sharding
One query is slow
Read EXPLAIN ANALYZE, add or fix an index
A missing index looks exactly like a table that is too big
Reads overwhelm the primary
Add a read replica or a cache
Spreads reads without splitting data, though replicas can lag
CPU, RAM or disk is exhausted
Move to a bigger machine
A larger box is a deploy, not a redesign
Old rows dominate the table
Archive them, or range-partition by date
Retention is the cheapest win partitioning offers
Writes exceed what one server sustains
Shard, after the steps above are spent
Only extra servers add write capacity
On the single-VPS setups I usually deploy, a Postgres box with Redis in front and a replica for backups goes a long way. For tables like an order log, the first step I would reach for is partitioning by month, not sharding.
When should you partition, and when is sharding unavoidable?
Use this checklist as a decision rule. Answer in order and stop at the first yes.
Is a query slow because of a missing index or a bad plan? Fix the index first.
Is the table large and time-ordered, with queries that filter on date and old data you want to drop? Range-partition on that date.
Is the table large but every query already filters on one key such as tenant, and a single partition would still fit on one server? Hash or list partitioning can help with maintenance, still on one machine.
Is the working set or the write rate larger than the biggest server you can reasonably run, even after indexes, replicas and partitioning? Sharding is now unavoidable.
Do you need data residency, with a region's data pinned to servers in that region? That can justify sharding by region even at modest size.
These two techniques also combine. Each shard can partition its own tables by date, which is how large multi-tenant systems keep retention cheap inside every shard. The order matters, though: partition first because it is reversible, and shard only when a measured limit forces you to.
The rule to carry: partitioning is a table-layout decision inside one server, sharding is a system-architecture decision across servers. Reach for partitioning when scan size or retention hurts, exhaust indexes, replicas and a bigger box next, and shard only when one machine cannot carry the writes. Choosing a shard key and a resharding plan is the real work, and no query planner will do it for you.