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.

9 min read
On this page 8 sections
  1. Partitioning vs sharding
  2. Postgres table partitioning
  3. Partitioning by date for test attempts and logs
  4. Sharding and shard keys
  5. The costs of sharding
  6. Do you need to shard?
  7. Key takeaways
  8. Frequently asked questions

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

AspectPartitioning (in one Postgres server)Sharding (across servers)
Where the data livesChild tables of one parent, on one serverSeparate databases, usually on separate machines
What it solvesHuge tables, slow maintenance, expensive deletes of old dataWrite load or data size beyond what one server can handle; isolating tenants
Application changesNone for most queries; the planner skips irrelevant partitionsEvery query must be routed to the right shard by a shard key
Transactions and joinsWork as normalEasy within one shard, hard across shards
Operational costLow: create future partitions, drop old onesHigh: 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:30 offset, 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_at too. "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:

StrategyHow rows are placedStrengthWeakness
HashA hash of the key picks the shardEven spreadRange queries touch every shard; adding shards means moving data
RangeKey ranges map to shardsSimple to split and reason aboutNew data can pile onto one "hot" shard
DirectoryA lookup table maps each key to a shardCan move individual tenantsAn extra lookup, and the directory must stay highly available
By tenantAll of one customer's data on one shardNatural for multi-tenant software; most queries stay on one shardOne 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?

SymptomTry firstSharding starts to make sense when
Slow queries on big tablesBetter database indexing, query fixes, partitioningNever for this alone
Heavy read loadCaching and read replicasRarely, since reads scale out without sharding
Primary saturated by writesBatching writes, a bigger server, moving event and log streams to other storesThe largest practical server still can't keep up after tuning
Database too big to restore within your RTOPartitioning and archiving cold dataA single-node restore still takes longer than the business can accept
One tenant's load hurting othersPer-tenant rate limits and separate queuesSome 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.

Share this article

Looking for something else?

Talk to Us