Sharding vs partitioning: how databases split data
Partition inside one server before you shard across many: Postgres declarative partitioning by date, shard keys for multi-institute platforms, the costs of sharding and when it pays.
On this page 8 sections
Partitioning splits one large table into smaller pieces inside the same database server; sharding splits data across several database servers, each holding a subset of the rows. Partitioning is a maintenance and performance tool you can adopt table by table with PostgreSQL's built-in declarative partitioning. Sharding changes the architecture, touching queries, transactions, migrations and backups, which is why most platforms should index, cache, add read replicas and partition long before they shard.
Partitioning vs sharding
| Aspect | Partitioning (in one Postgres server) | Sharding (across servers) |
|---|---|---|
| Where the data lives | Child tables of one parent, on one server | Separate databases, usually on separate machines |
| What it solves | Huge tables, slow maintenance, expensive deletes of old data | Write load or data size beyond what one server can handle; isolating tenants |
| Application changes | None for most queries; the planner skips irrelevant partitions | Every query must be routed to the right shard by a shard key |
| Transactions and joins | Work as normal | Easy within one shard, hard across shards |
| Operational cost | Low: create future partitions, drop old ones | High: routing, rebalancing, per-shard migrations and backups |
The terms overlap in casual use. Sharding is horizontal partitioning taken across machines: rows are divided up, and each machine holds some of them. Vertical partitioning, by contrast, splits columns or whole tables apart, such as moving bulky analytics events into a different database from enrolments and payments.
Postgres table partitioning
PostgreSQL has had declarative partitioning since version 10, with range and list partitioning; version 11 added hash partitioning and a default partition. A partitioned table is a parent with no rows of its own. Each child partition holds the rows whose partition key falls in its range or list, and Postgres routes inserts to the right child automatically. The PostgreSQL partitioning documentation sets out what it buys you:
- Faster queries on recent data. When a query's WHERE clause matches the partition key, the planner prunes partitions it doesn't need, and the busy partitions' indexes are more likely to fit in memory.
- Cheap removal of old data. Dropping or detaching a partition is far faster than a bulk
DELETE, and avoids the vacuum work a bulk delete leaves behind. - Cheaper storage for cold data. Rarely used partitions can live on slower, cheaper storage.
The same documentation offers a rule of thumb: partitioning starts to pay off when a table is larger than the database server's memory. It also warns against too many partitions: the planner copes with a few thousand when queries prune most of them away, but planning time and memory grow as more partitions survive pruning.
Know the limitations before you design. A primary key or unique constraint on a partitioned table must include every partition-key column, because each partition can only enforce uniqueness within itself. And you pick the partition key once: choose the column that appears most often in your WHERE clauses.
Partitioning by date for test attempts and logs
Append-mostly data that you query by time and eventually delete is the classic case: answer submissions, video playback events, login audits and notification logs. An illustrative sizing: 3 lakh active students, each taking eight tests a month of 100 questions, produce 24 crore answer rows a month. At a rough 130 bytes per row including two indexes, that's about 30 GB a month, and more than half a terabyte after 18 months. Deleting a month's worth of old rows from one giant table is painful; dropping a monthly partition is almost instant.
CREATE TABLE attempt_answers (
id bigserial,
attempt_id bigint NOT NULL,
question_id bigint NOT NULL,
answered_at timestamptz NOT NULL,
choice smallint,
PRIMARY KEY (id, answered_at)
) PARTITION BY RANGE (answered_at);
CREATE TABLE attempt_answers_2026_10 PARTITION OF attempt_answers
FOR VALUES FROM ('2026-10-01 00:00+05:30') TO ('2026-11-01 00:00+05:30');
CREATE INDEX ON attempt_answers (attempt_id, answered_at);
Points worth noticing:
- The primary key includes
answered_at, as the uniqueness rule requires. - The bounds carry an explicit
+05:30offset, so each month starts at midnight in India rather than at whatever time zone the session happens to use. A range includes its lower bound and excludes its upper bound. - An index created on the parent is created on every partition automatically, including future ones.
- Queries should filter on
answered_attoo. "All answers for attempt 812" without a date range checks every partition's index; add the attempt's time window and Postgres prunes to one or two partitions.
Create partitions ahead of time, never at the moment the first row arrives. A scheduled job can create next month's partition, or the pg_partman extension can create future partitions and drop old ones under a retention policy. To retire a month without blocking queries, use ALTER TABLE ... DETACH PARTITION ... CONCURRENTLY (PostgreSQL 14 and later), archive the detached table to object storage, then drop it. Two cautions: CONCURRENTLY can't run inside a transaction block, and it isn't allowed while the table has a default partition, so think twice before adding one.
In Django, the ORM can query a partitioned table like any other, but migrations won't create one for you; use RunSQL with state_operations so Django's model state still matches. Django 5.2 added CompositePrimaryKey, which can model a primary key such as (id, answered_at), but Django can't migrate an existing table to a composite key, and other models can't point a ForeignKey at it. That's rarely a problem for a leaf table like this one. If you're partitioning an existing large table, move data across in batches using the approach in our guide to zero-downtime migrations.
Sharding and shard keys
A shard key decides which server holds each row, and it's the most consequential choice in a sharded design. A good one is present in almost every query, keeps each transaction on one shard, spreads load evenly and never changes for a given row. Common strategies:
| Strategy | How rows are placed | Strength | Weakness |
|---|---|---|---|
| Hash | A hash of the key picks the shard | Even spread | Range queries touch every shard; adding shards means moving data |
| Range | Key ranges map to shards | Simple to split and reason about | New data can pile onto one "hot" shard |
| Directory | A lookup table maps each key to a shard | Can move individual tenants | An extra lookup, and the directory must stay highly available |
| By tenant | All of one customer's data on one shard | Natural for multi-tenant software; most queries stay on one shard | One very large tenant can outgrow its shard |
If you run a platform that serves many coaching institutes, the institute is the natural tenant key: a student's enrolments, attempts, payments and progress all belong to one institute, so nearly every query and transaction stays on one shard. The Citus documentation gives the same advice for multi-tenant applications: distribute tables by the tenant ID, co-locate related tables on that same column so joins and foreign keys keep working, and turn small shared tables into reference tables copied to every node.
Does Postgres shard on its own? Core PostgreSQL has no built-in automatic sharding. The main routes are the Citus extension, which distributes tables and queries across nodes; application-level sharding, such as Django database routers choosing a database per institute; and, for simpler cases, partitions defined as foreign tables on other servers through postgres_fdw, which comes with limitations such as no unique indexes on the parent.
The costs of sharding
- Cross-shard queries. An all-India leaderboard across institutes, or a platform-wide revenue report, must gather results from every shard.
- Cross-shard transactions. Anything touching two shards needs two-phase commit or a redesign so it doesn't.
- Global uniqueness. A unique email or phone number across all shards can't be enforced by one index.
- Rebalancing. Moving a growing tenant to a new shard is a live data migration of its own.
- Operations times N. Every schema migration, backup, restore test and upgrade runs on every shard. Restoring several shards to one consistent moment after an incident is much harder than restoring one database; see point-in-time recovery.
- Analytics move out. Reporting usually needs a separate warehouse fed from all shards.
Do you need to shard?
| Symptom | Try first | Sharding starts to make sense when |
|---|---|---|
| Slow queries on big tables | Better database indexing, query fixes, partitioning | Never for this alone |
| Heavy read load | Caching and read replicas | Rarely, since reads scale out without sharding |
| Primary saturated by writes | Batching writes, a bigger server, moving event and log streams to other stores | The largest practical server still can't keep up after tuning |
| Database too big to restore within your RTO | Partitioning and archiving cold data | A single-node restore still takes longer than the business can accept |
| One tenant's load hurting others | Per-tenant rate limits and separate queues | Some tenants need dedicated capacity or isolation |
Scaling a single database server up is usually the cheaper step for a long time; our guide to horizontal vs vertical scaling explains why databases tend to scale vertically first. When you do shard, design the shard key around your tenants and your most frequent queries, and expect the operational load to grow with every shard.
Key takeaways
- Partitioning splits a table within one server; sharding splits data across servers.
- Postgres declarative partitioning suits time-based data such as attempts and logs: fast pruning, near-instant retention and few application changes.
- Partition keys and primary keys are linked, and bounds should state their time zone explicitly.
- If you shard a multi-institute platform, shard by institute and keep related tables co-located.
- Index, cache, replicate, partition and scale up first; shard when a single primary truly can't keep up.
Frequently asked questions
Is sharding and partitioning the same?
Not quite. Partitioning is the general idea of splitting data into parts; in PostgreSQL it usually means dividing a table into child tables on the same server. Sharding is partitioning across separate database servers, so each server holds only some of the rows. Partitioning needs no application changes, while sharding requires routing every query to the right server.
Is sharding horizontal scaling?
Yes. Sharding is a way to scale a database horizontally: instead of buying a bigger server, you add servers and give each a share of the data and the write load. Read replicas also scale out, but only for reads, since every replica holds a full copy and all writes still go to one primary. Sharding is how write capacity grows beyond one machine.
Does Postgres support sharding?
Core PostgreSQL has declarative partitioning but no built-in automatic sharding across servers. You can shard with the Citus extension, which distributes tables and queries across nodes; in the application, by routing each tenant to its own database; or in limited cases by making partitions foreign tables on other servers with postgres_fdw. Each approach adds operational work.
When to use sharding?
Use sharding when one primary database can no longer handle your write load or data size even after indexing, query tuning, caching, read replicas, partitioning and moving to the largest practical server, or when some tenants need isolated capacity. For most learning platforms, that point comes much later than expected, and partitioning large time-based tables solves the more common problems first.