← Back to the section

When a table grows to tens of gigabytes, problems begin: queries slow down, VACUUM runs for hours, and deleting old data turns into a multi-hour operation. Let's look at two tools — partitioning and sharding: how they differ and when to use each.

PARTITION BY RANGE (created_at) product_2026_q1 2026-01-01 → 2026-04-01 product_2026_q2 2026-04-01 → 2026-07-01 product_default anything outside q1 and q2 INSERTcreated_at = 2026-05-14 product_2026_q22026-04-01 → 2026-07-01+ rowkey 2026-05-14 falls into the q2 range — the row lands there SELECTcreated_at >= 2026-04-01 product_2026_q22026-04-01 → 2026-07-01product_defaultanything outside q1 and q2prunedscannedscannedcondition open on the right: q1 is pruned, q2 and DEFAULT are scanned

One logical table and three physical pieces. On insert PostgreSQL looks at the key value and puts the row into the partition whose range holds it; a date outside every range goes to DEFAULT. On read the planner throws away partitions whose range cannot match the condition — but it cannot throw away DEFAULT while the condition is open on the right: anything at all may be sitting there.

How PostgreSQL stores data and why a large table is a pain

When a query hits a table, PostgreSQL works with pages (8 KB each). An index helps find the needed pages quickly. But if the table takes up 50–100 GB, even the index becomes large — it doesn't fit entirely in RAM (shared_buffers), and reading turns into constant disk access.

Another pain point is cleaning up outdated row versions. PostgreSQL doesn't delete old versions right away; that's done by the background VACUUM process. On a large table it can run around the clock, unable to keep up with the flow of changes.

There are two strategies for solving these problems:

  • Partitioning — cut one large table into several physical pieces within a single database on a single server.
  • Sharding — spread data across several physical servers.

The first is optimization within a single machine, the second is scaling across several.

How to find out how much space a table takes

Before optimizing anything, you need to measure. PostgreSQL has handy functions:

-- Size of the entire database
SELECT pg_size_pretty(pg_database_size('shop'));
-- → 42 GB

-- Size of the table together with indexes and TOAST
SELECT pg_size_pretty(pg_total_relation_size('product'));
-- → 12 GB

-- Just the rows themselves (without indexes)
SELECT pg_size_pretty(pg_relation_size('product'));
-- → 7 GB

-- Indexes only
SELECT pg_size_pretty(pg_indexes_size('product'));
-- → 5 GB

TOAST is separate storage for long values (texts, JSON). If pg_total_relation_size is much larger than the sum of pg_relation_size and pg_indexes_size, it means the table stores a lot of long fields.

To see all tables sorted by size:

live example

SELECT
    relname,
    pg_size_pretty(pg_total_relation_size(schemaname || '.' || relname)) AS total,
    pg_size_pretty(pg_relation_size(schemaname || '.' || relname))       AS table_only,
    pg_size_pretty(pg_indexes_size(schemaname || '.' || relname))        AS indexes,
    n_live_tup   AS live_rows,
    n_dead_tup   AS dead_rows
FROM pg_stat_user_tables
ORDER BY pg_total_relation_size(schemaname || '.' || relname) DESC
LIMIT 20;
Run

Running examples is part of paid access. There the same code runs inside the article: editor, run and check next to the paragraph. Three free days →

If dead_rows is comparable to live_rows — VACUUM isn't keeping up with deletions. If indexes is larger than table_only — you may have too many indexes.

Partitioning

Partitioning is when one logical table is physically cut into several parts called partitions. The application still writes to and reads from the product table, and PostgreSQL itself figures out which physical partition to put a row into and which one to read from.

This is called partition pruning — the planner "cuts off" unnecessary partitions and reads only those where the needed rows might be.

Partitioning solves the problems of a large table on a single server:

  • indexes on each partition are smaller and fit into memory better;
  • VACUUM works on partitions separately — faster and lighter;
  • dropping an entire partition (DROP TABLE product_2023) is an instant operation instead of a multi-hour DELETE.

RANGE — by range

The most popular option is to cut by dates. Logs, orders, transactions — anything that accumulates over time fits well into RANGE partitions.

CREATE TABLE product (
    id          BIGSERIAL,
    category_id BIGINT,
    price       NUMERIC(10, 2) NOT NULL,
    name        TEXT NOT NULL,
    created_at  TIMESTAMPTZ NOT NULL,
    PRIMARY KEY (id, created_at)   -- the partition key must be part of the PRIMARY KEY
) PARTITION BY RANGE (created_at);

CREATE TABLE product_2026_q1 PARTITION OF product
    FOR VALUES FROM ('2026-01-01') TO ('2026-04-01');

CREATE TABLE product_2026_q2 PARTITION OF product
    FOR VALUES FROM ('2026-04-01') TO ('2026-07-01');

The query SELECT * FROM product WHERE created_at >= '2026-04-01' reads only product_2026_q2. The other partitions aren't even opened.

LIST — by specific values

When data is divided by a finite set of values — for example, by country, region, or status.

CREATE TABLE product (
    id          BIGSERIAL,
    category_id BIGINT NOT NULL,
    price       NUMERIC(10, 2) NOT NULL,
    name        TEXT NOT NULL,
    PRIMARY KEY (id, category_id)
) PARTITION BY LIST (category_id);

CREATE TABLE product_sweets PARTITION OF product FOR VALUES IN (1);
CREATE TABLE product_meat   PARTITION OF product FOR VALUES IN (2);
CREATE TABLE product_dairy  PARTITION OF product FOR VALUES IN (3);
CREATE TABLE product_other  PARTITION OF product DEFAULT;

The query WHERE category_id = 1 reads only product_sweets. The planner skips the rest.

HASH — even distribution

When there's no natural scale and you just need to evenly distribute the load across partitions. PostgreSQL takes a hash of the key value and distributes rows by the formula hash(id) % N.

CREATE TABLE product (
    id          BIGSERIAL,
    category_id BIGINT,
    price       NUMERIC(10, 2) NOT NULL,
    name        TEXT NOT NULL,
    PRIMARY KEY (id)
) PARTITION BY HASH (id);

CREATE TABLE product_p0 PARTITION OF product FOR VALUES WITH (MODULUS 8, REMAINDER 0);
CREATE TABLE product_p1 PARTITION OF product FOR VALUES WITH (MODULUS 8, REMAINDER 1);
-- ... product_p7

The query WHERE id = 12345 reads only one partition. But a query without a filter on id still reads all eight.

One more thing: hash partitioning has no DEFAULT partition — a key value always lands in one of the remainders, so "leftover" rows simply don't exist here. DEFAULT is only for RANGE and LIST.

The DEFAULT partition

A typical failure with RANGE partitioning: a new quarter arrives, the partition hasn't been created yet, and the INSERT fails with the error no partition of relation "product" found for row. The row is not stored at all, and the application spews errors.

The safeguard is the DEFAULT partition: it catches all rows that didn't land in any declared partition.

CREATE TABLE product_default PARTITION OF product DEFAULT;

Now the INSERT won't fail. But DEFAULT is insurance, not a place to store data. If a lot of rows accumulate there, then when you try to ATTACH PARTITION a new range, PostgreSQL will first scan the DEFAULT partition under a lock, and that can take a long time.

A second quirk: DEFAULT is not always pruned. Its range is "everything else", so the planner can drop it from the plan only when the condition fits entirely inside the declared ranges. The condition created_at >= '2026-04-01' is open on the right: a row from 2027 may be sitting in DEFAULT, and it has to be read.

Working with DEFAULT correctly:

  1. Create new partitions in advance — several periods ahead (via cron or pg_partman).
  2. Set up an alert if anything shows up in DEFAULT — that means the auto-creation broke.
  3. Before each ATTACH, make sure DEFAULT is empty.

Where a row lands and what a query reads

The whole mechanic is a function of the key value: on write PostgreSQL picks the physical table from created_at, on read it uses the condition to decide which tables to open at all. The same logic without a database, in plain Java:

live example

import java.time.LocalDate;

public class PartitionRouting {
    record Partition(String name, LocalDate from, LocalDate to) {
        boolean holds(LocalDate key) {
            return !key.isBefore(from) && key.isBefore(to);
        }
    }

    static final Partition[] PARTS = {
            new Partition("product_2026_q1", LocalDate.parse("2026-01-01"), LocalDate.parse("2026-04-01")),
            new Partition("product_2026_q2", LocalDate.parse("2026-04-01"), LocalDate.parse("2026-07-01"))};

    public static void main(String[] args) {
        for (String date : new String[]{"2026-02-11", "2026-05-14", "2026-11-02"}) {
            LocalDate key = LocalDate.parse(date);
            String target = "product_default";
            for (Partition p : PARTS) {
                if (p.holds(key)) {
                    target = p.name();
                }
            }
            System.out.println("INSERT created_at=" + date + " -> " + target);
        }

        LocalDate from = LocalDate.parse("2026-04-01");
        System.out.println("SELECT ... WHERE created_at >= " + from);
        for (Partition p : PARTS) {
            System.out.println("  " + p.name() + (p.to().isAfter(from) ? " scanned" : " pruned"));
        }
        System.out.println("  product_default scanned: no bound on the right");
    }
}
Run

Running examples is part of paid access. There the same code runs inside the article: editor, run and check next to the paragraph. Three free days →

The November row found no range of its own and landed in product_default; the condition created_at >= '2026-04-01' pruned product_2026_q1 but not DEFAULT.

How to choose a partition key

Three mandatory conditions:

  1. The key is present in most queries — otherwise the planner can't cut off the extra partitions and will read all of them.
  2. Data is distributed evenly — if 95% of rows land in one partition, there's no point.
  3. The key almost never changes — changing the key physically moves the row from one partition to another, and on hot updates that's expensive.

One more restriction: the partition key must be part of the PRIMARY KEY and every UNIQUE index. That's exactly why the examples above use PRIMARY KEY (id, created_at), not just PRIMARY KEY (id).

Sharding

Sharding is when data is spread across several physical servers. Each server (shard) stores only part of the data. The application or a proxy layer knows which shard to go to for each specific query.

Sharding is needed when a single server can no longer cope: it's short on RAM, IOPS, or disk throughput. Partitioning won't solve this — the data is still physically on one machine.

Sharding strategies

By hash — the simplest. shard = hash(id) % N. Data is distributed evenly. The problem: when adding a new server, you need to move a significant portion of the data. This is solved with consistent hashing, which minimizes the number of rows moved.

By range — data with id from 0 to 10M on one shard, from 10M to 20M on another. It's easy to add new servers, but the last range is always "hot" — all new records go there.

By tenant — each client (tenant) gets its own shard. Convenient for SaaS: a large client gets a dedicated server, and data isolation is simple. Requires a routing table.

By geography — data for users in Europe, Asia, and the US is stored on servers in the corresponding regions. This reduces latency and helps with requirements to store data in specific countries.

How to choose a shard key

Everything that applies to the partition key, plus two additional requirements:

  • Related data should live on the same shard — if product is sharded by category_id and category is also sharded by category_id, then a JOIN between them stays local within the shard. Otherwise every JOIN turns into a distributed query.
  • Transactions should fit within a single shard — distributed transactions across several servers are slow and cope poorly with network failures.

The cost of sharding

Sharding adds complexity that partitioning doesn't have:

  • Distributed transactions — kept out by design: one operation = one shard.
  • JOINs across shards — almost always a sign of a wrong key choice: you have to denormalize data or synchronize it asynchronously.
  • Global uniqueness — BIGSERIAL doesn't work: counters on different shards diverge. You need UUIDs or special identifier generators.
  • Adding a new server — requires moving data. This is slow and needs special infrastructure; for PostgreSQL that is Citus.

Partitioning vs sharding

PartitioningSharding
Where the data isDifferent tables in one databaseDifferent servers
When it's neededTable > 50–100 GB, index and VACUUM slowdownsA single server can no longer handle it
Transparent to the applicationYes — one logical tableNo — routing is needed
Transactions across segmentsOrdinary local onesDistributed — expensive
JOINs across segmentsCheapExpensive
Changing the keyExpensive but possiblePractically impossible
Toolspg_partman, scriptsCitus, manual work

Practical rule: partition first, shard only when you've hit the ceiling of a single server. Partitioning solves 80% of large-table problems and is incomparably easier to manage.

A useful trick for those planning to shard in the future: partition by the same key you later want to shard by. Then the move to sharding is a matter of transferring ready-made partitions to other servers, without rewriting queries.

In short

  • Partitioning — one table, several physical pieces on a single server. Sharding — data on several servers.
  • RANGE — by range (dates, numbers). LIST — by specific values. HASH — even distribution.
  • The partition key must be part of the PRIMARY KEY and all UNIQUE indexes, and must appear in the condition of most queries — otherwise the planner reads every partition.
  • The DEFAULT partition is insurance against failing INSERTs and is normally empty; the planner cannot always prune it. New partitions are created in advance (via cron or pg_partman).
  • Sharding makes JOINs across shards and transactions across several shards expensive — this is avoided by design.
  • A global BIGSERIAL doesn't work with sharding — you need UUIDs or special generators.