What Is Database Partitioning?

By Oleksandr Andrushchenko — Published on
Category: Databases & Data
0 Likes
0 Dislikes
What Is Database Partitioning?
What Is Database Partitioning?

Database partitioning divides a large table into smaller physical pieces called partitions while preserving a single logical table from the application's perspective.

Partitioning can improve query performance through partition pruning, simplify retention and maintenance, reduce the amount of data touched by some operations, and make very large tables easier to operate. It does not automatically distribute the database across multiple servers, and it should not be confused with database sharding.

Table of Contents

Why Database Partitioning Exists

Consider an application that stores several billion events in one table:

events
--------------------------------
id
account_id
event_type
payload
created_at

5,000,000,000+ rows

The database can technically store the table, but operating it becomes increasingly difficult.

Problems may include:

  • queries touching large amounts of irrelevant data;
  • very large indexes;
  • slow index maintenance;
  • expensive deletion of historical records;
  • long-running maintenance operations;
  • large vacuum or cleanup workloads;
  • increasing backup and operational complexity.

Suppose most application queries access recent events:

SELECT id, event_type, payload, created_at
FROM events
WHERE account_id = $1
  AND created_at >= $2
  AND created_at < $3
ORDER BY created_at DESC;

If the table contains five years of data but the query requests only one month, treating the entire dataset as one physical unit may be unnecessary.

The table can instead be divided by time:

events
  ├── events_2026_07
  ├── events_2026_08
  ├── events_2026_09
  ├── events_2026_10
  └── ...

The application still queries events, while the database determines which physical partitions contain the requested rows.

How Database Partitioning Works

A partitioned table remains one logical table but stores rows in separate partitions according to a partitioning rule.

Logical Table: events
          │
          ├──→ Partition A
          ├──→ Partition B
          ├──→ Partition C
          └──→ Partition D

When a row is inserted, its partition key determines where it belongs.

For a table partitioned monthly by created_at:

2026-07-14 → events_2026_07
2026-08-03 → events_2026_08
2026-09-29 → events_2026_09

In PostgreSQL, a simplified definition could look like:

CREATE TABLE events (
    id BIGINT NOT NULL,
    account_id BIGINT NOT NULL,
    event_type TEXT NOT NULL,
    payload JSONB,
    created_at TIMESTAMPTZ NOT NULL
) PARTITION BY RANGE (created_at);

Monthly partitions can then be created:

CREATE TABLE events_2026_09
PARTITION OF events
FOR VALUES FROM ('2026-09-01')
TO ('2026-10-01');

CREATE TABLE events_2026_10
PARTITION OF events
FOR VALUES FROM ('2026-10-01')
TO ('2026-11-01');

The application continues querying the parent table:

SELECT *
FROM events
WHERE created_at >= '2026-09-01'
  AND created_at < '2026-10-01';

The database can determine that only the September partition is relevant.

A Guide To Database Sharding
A Guide To Database Sharding

Choosing a Partition Key

The partition key determines how rows are assigned to partitions.

This is the most important partitioning decision because it affects query pruning, data distribution, maintenance, and future growth.

Common partition keys include:

  • timestamp;
  • tenant or account ID;
  • region;
  • status or category;
  • numeric ID;
  • a hash of an identifier.

A useful partition key should usually appear naturally in important queries.

For example, if nearly every event query contains:

WHERE created_at >= ...
  AND created_at < ...

then created_at is a strong candidate for range partitioning.

If queries instead look like:

WHERE account_id = $1

partitioning by account may provide better pruning.

A partition key selected only because it distributes rows evenly can still perform poorly if application queries cannot use it.

Partitioning Strategies

Different partitioning strategies suit different access patterns. Three common approaches are range, list, and hash partitioning.

Range Partitioning

Range partitioning assigns rows according to ordered value ranges.

created_at

Jan 2026 → Partition 1
Feb 2026 → Partition 2
Mar 2026 → Partition 3
Apr 2026 → Partition 4

Time-series data is a natural use case.

Another example could partition numeric IDs:

1 - 1,000,000
        → Partition A

1,000,001 - 2,000,000
        → Partition B

2,000,001 - 3,000,000
        → Partition C

Range partitioning works well when queries frequently request contiguous ranges.

The main risk is uneven activity. If all new writes go to the newest time partition, that partition receives nearly all write traffic.

List Partitioning

List partitioning maps specific values to partitions.

For example, orders could be partitioned by region:

US, CA
   → north_america

DE, FR, IT, ES
   → europe

JP, KR
   → asia

A simplified PostgreSQL definition could be:

CREATE TABLE orders (
    id BIGINT NOT NULL,
    region TEXT NOT NULL,
    total_cents BIGINT NOT NULL
) PARTITION BY LIST (region);

List partitioning is useful when data naturally belongs to a small number of known categories.

It becomes harder to manage when values change frequently or when one category becomes much larger than the others.

Hash Partitioning

Hash partitioning applies a hash function to a key and distributes rows among a fixed number of partitions.

hash(account_id) % 4

0 → Partition 0
1 → Partition 1
2 → Partition 2
3 → Partition 3

Unlike time-based range partitioning, writes can be spread more evenly across partitions.

For example:

account 152 → hash → Partition 3
account 981 → hash → Partition 0
account 427 → hash → Partition 2

Hash partitioning is useful when equality lookups dominate and reasonably even distribution matters.

It is less convenient for range queries because adjacent values do not necessarily live in the same partition.

Partition Pruning

One of the main performance benefits of partitioning is partition pruning.

Partition pruning allows the database optimizer to exclude partitions that cannot contain rows matching a query.

Consider monthly partitions:

events_2026_06
events_2026_07
events_2026_08
events_2026_09

A query asks only for September:

SELECT *
FROM events
WHERE created_at >= '2026-09-01'
  AND created_at < '2026-10-01';

The optimizer can reduce the operation to:

events_2026_06 → skip
events_2026_07 → skip
events_2026_08 → skip
events_2026_09 → scan

Instead of considering the entire dataset, the database works only with the relevant partition.

Partitioning becomes much less useful when important queries do not filter by the partition key.

For example:

SELECT *
FROM events
WHERE event_type = 'payment_failed';

If the table is partitioned by created_at and the query contains no time condition, the database may need to search every partition.

Partition 1 ─┐
Partition 2 ─┤
Partition 3 ─┼──→ Query
Partition 4 ─┤
Partition 5 ─┘

This is why partitioning should follow actual access patterns rather than being added merely because a table is large.

Partitioning and Indexes

Partitioning does not replace indexing.

Partition pruning answers:

Which partitions need to be searched?

An index answers:

How can matching rows be found efficiently
inside those partitions?

Consider:

SELECT id, event_type, created_at
FROM events
WHERE account_id = 4821
  AND created_at >= '2026-09-01'
  AND created_at < '2026-10-01'
ORDER BY created_at DESC;

Time partitioning can prune the query to the September partition.

An index such as:

CREATE INDEX ON events (account_id, created_at DESC);

can then help locate that account's rows efficiently within the selected partition.

The combination becomes:

Partition pruning
      ↓
September partition only
      ↓
Index lookup
      ↓
Account 4821 rows

Partitioning and indexing solve different layers of the same query problem.

Partitioning vs Sharding

Partitioning and sharding both divide data, but the operational boundary is different.

Database Partitioning Database Sharding
Usually one database cluster Multiple independent database nodes or clusters
One logical table Dataset distributed across shards
Database usually handles row placement Application or middleware may handle routing
Cross-partition queries remain inside the database Cross-shard queries may require multiple nodes
Normal local transactions remain available Cross-shard transactions become harder
Primarily improves manageability and selective access Can increase total storage and write capacity

A partitioned PostgreSQL table may contain 100 partitions while all of them still live inside the same database cluster.

PostgreSQL
   │
   └── events
       ├── partition_1
       ├── partition_2
       └── partition_3

Sharding moves data across independent database systems:

Application
   │
   ├──→ Database Shard A
   ├──→ Database Shard B
   └──→ Database Shard C

Partitioning therefore does not automatically solve a database that has exhausted the CPU, memory, connection, storage I/O, or write throughput of one server.

Sharding addresses a different scaling boundary and introduces substantially more routing and coordination complexity. Those trade-offs are covered in What Is Database Sharding?.

Hot Partitions

A hot partition receives significantly more traffic than other partitions.

Time partitioning naturally creates this possibility:

January   → historical
February  → historical
March     → historical
...
September → all current writes

Historical partitions may receive almost no writes while the current partition handles the entire ingestion workload.

Tenant partitioning can create a similar problem:

Tenant A → 500 requests/sec
Tenant B → 300 requests/sec
Tenant C → 40,000 requests/sec

If Tenant C maps to one partition, that partition may become overloaded even though the total dataset is well distributed.

A good partitioning strategy considers both data size and traffic distribution.

Hash partitioning can distribute unrelated keys more evenly, but it does not automatically fix every hotspot. A single extremely popular key can still dominate the partition that owns it.

Partition Size and Count

Partitions should be large enough to justify their existence but small enough to provide useful pruning and operational boundaries.

Consider a table containing ten years of event data.

One partition per year produces:

10 partitions

One partition per month produces:

120 partitions

One partition per day produces approximately:

3,650 partitions

One partition per hour produces approximately:

87,600 partitions

More partitions are not automatically better.

Very large partitions reduce pruning precision and make maintenance units large. Excessively small partitions increase metadata, planning, schema-management, and operational overhead.

The appropriate interval depends on:

  • data ingestion rate;
  • query time ranges;
  • retention policy;
  • maintenance requirements;
  • database engine behavior;
  • expected partition size.

A system ingesting billions of events per day may need much smaller time partitions than an application storing a few million rows per month.

Retention and Data Lifecycle

Time partitioning is particularly useful when data has a retention period.

Suppose event data must be retained for 12 months.

Without partitioning, removing old data might require:

DELETE FROM events
WHERE created_at < '2025-10-01';

Deleting hundreds of millions of rows can generate substantial transaction logs, locking pressure, vacuum work, replication traffic, and I/O.

With monthly partitions:

events_2025_07
events_2025_08
events_2025_09
events_2025_10
...
events_2026_09

old data can often be removed at the partition level rather than row by row.

DROP TABLE events_2025_09;

Production systems may detach, archive, verify, and then delete an old partition instead of dropping it immediately.

This makes partition boundaries useful operational lifecycle boundaries, not only query-performance boundaries.

Repartitioning

Partitioning decisions can become incorrect as workloads change.

Suppose an events table initially uses yearly partitions:

2024
2025
2026

Traffic grows until the 2026 partition contains several billion rows.

Monthly partitions may now be more appropriate:

2026-01
2026-02
2026-03
...
2026-12

Changing the partitioning scheme can require moving large amounts of data.

A safe migration may involve:

Create new partitioned structure
        ↓
Backfill historical data in batches
        ↓
Keep new writes synchronized
        ↓
Validate counts and checksums
        ↓
Switch reads
        ↓
Switch final write path
        ↓
Retire old structure

Repartitioning billions of rows is an operational project, not a trivial schema migration.

This is another reason to select partition keys from stable access patterns and expected data lifecycle rather than short-term convenience.

Production Design Example

Consider a SaaS platform storing customer audit events.

The system ingests:

150 million events/day
≈ 4.5 billion events/month

Most customer queries request between one hour and seven days of history:

SELECT id, event_type, actor_id, created_at
FROM audit_events
WHERE account_id = $1
  AND created_at >= $2
  AND created_at < $3
ORDER BY created_at DESC
LIMIT 500;

The retention requirement is 13 months.

A single unpartitioned table makes retention and maintenance increasingly expensive.

The table is therefore range-partitioned by day:

audit_events
   │
   ├── 2026-09-27
   ├── 2026-09-28
   ├── 2026-09-29
   ├── 2026-09-30
   └── ...

The schema uses:

CREATE TABLE audit_events (
    id BIGINT NOT NULL,
    account_id BIGINT NOT NULL,
    actor_id BIGINT,
    event_type TEXT NOT NULL,
    payload JSONB,
    created_at TIMESTAMPTZ NOT NULL
) PARTITION BY RANGE (created_at);

Each partition receives an index supporting the dominant query pattern:

CREATE INDEX ON audit_events
    (account_id, created_at DESC);

A customer requesting two days of data produces:

13 months of partitions
        ↓
Partition pruning
        ↓
2 relevant daily partitions
        ↓
(account_id, created_at) indexes
        ↓
Return newest 500 rows

The system automatically creates future partitions before they are needed:

Today: 2026-09-29

Already created:
2026-09-30
2026-10-01
2026-10-02

This avoids a production failure where a new row arrives but no partition exists for its timestamp.

The retention workflow runs separately:

Partition older than 13 months
        ↓
Detach partition
        ↓
Archive if required
        ↓
Verify archive
        ↓
Drop partition

This avoids enormous row-by-row delete jobs against the active table.

Important production metrics include:

  • rows and bytes per partition;
  • partition growth rate;
  • queries touching many partitions;
  • query planning time;
  • partition-pruning effectiveness;
  • index size per partition;
  • current-partition write throughput;
  • retention job duration;
  • missing future partitions;
  • database CPU and storage I/O.

As long as the database cluster still has enough aggregate write, CPU, memory, connection, and storage capacity, partitioning keeps the operational model relatively simple.

If the single cluster eventually becomes the bottleneck, horizontal database scaling may become necessary. Database Scaling Explained: Vertical vs Horizontal Scaling covers that progression.

Failure Scenarios

Partitioning introduces operational failure modes that should be planned explicitly.

The next partition was never created. New inserts can fail when rows no longer match an existing partition.

A query omits the partition key. The database may scan many or all partitions, creating unexpectedly expensive queries.

One partition grows much faster than expected. Maintenance and query performance become uneven.

A retention job removes the wrong partition. Large amounts of data can disappear almost instantly, so partition lifecycle operations require strong safeguards.

A new index is created incorrectly. Some partitions may have different indexing behavior or deployment state than others, depending on the database and migration process.

Repartitioning is interrupted. The migration needs checkpoints and validation so it can resume without duplicating or losing data.

The database server itself becomes saturated. Partitioning cannot provide additional physical CPU, memory, connection, or storage capacity simply by creating more partitions.

When to Use Database Partitioning

Database partitioning is useful when:

  • a table has grown very large;
  • queries naturally filter by time, tenant, region, or another partition key;
  • partition pruning can eliminate large amounts of irrelevant data;
  • historical data has a clear retention policy;
  • maintenance should operate on smaller units;
  • old data should be archived or removed efficiently;
  • indexes have become difficult to maintain as one physical structure;
  • one database cluster still has enough overall capacity.

Partitioning should not be introduced merely because a table contains millions of rows. Modern relational databases can handle large unpartitioned tables efficiently when indexes, queries, and maintenance are designed correctly.

A useful production progression is to optimize schemas and queries first, introduce partitioning when access patterns and lifecycle requirements justify it, and move to distributed data placement only when a single database reaches a measured capacity boundary.

That broader approach is covered in Database Best Practices for Scalable Applications.

Common Partitioning Mistakes

  • Partitioning without a measured problem. It adds operational complexity without guaranteed performance benefits.
  • Choosing a key unrelated to query patterns. Queries cannot prune partitions effectively.
  • Assuming partitioning replaces indexes. Queries still need efficient access inside selected partitions.
  • Creating too many tiny partitions. Planning, metadata, and maintenance overhead increase.
  • Creating partitions that are too large. Pruning and maintenance benefits become limited.
  • Forgetting future partitions. New inserts can fail at a partition boundary.
  • Ignoring hot partitions. Even data size does not guarantee even workload distribution.
  • Using time partitioning for queries without time filters. Many queries still touch every partition.
  • Assuming partitioning scales writes across servers. All partitions may still share the same database resources.
  • Confusing partitioning with sharding. The operational and transactional models are different.
  • Running huge row-by-row retention deletes. Time partitions can often make retention much cheaper.
  • Ignoring repartitioning cost. A poor initial key can become expensive to change after billions of rows exist.

Frequently Asked Questions

Partitioning is often discussed together with indexing, sharding, and database scaling, but each solves a different part of the problem.

Does Partitioning Improve Every Query?

No. Partitioning provides the largest query benefit when the database can prune irrelevant partitions.

A query that does not constrain the partition key may need to access many or all partitions and can sometimes perform worse than an equivalent query against a well-indexed unpartitioned table.

Does Partitioning Require Multiple Servers?

No. Table partitioning commonly divides data into multiple physical partitions inside the same database system.

Distributing data across independent database servers is generally associated with sharding or a distributed database architecture.

Can a Table Have Too Many Partitions?

Yes. Every database engine has operational and planning costs associated with large partition counts.

The correct number depends on the engine, data volume, access patterns, retention policy, and maintenance requirements. Partition granularity should solve a concrete operational or query problem rather than maximize the number of partitions.

Should Partitioning Be Used Before Sharding?

Often, yes, when the problem is large-table manageability, retention, or query pruning and one database still has enough physical capacity.

Partitioning keeps data inside one database environment and avoids much of the routing, transaction, migration, and operational complexity introduced by sharding. A detailed production implementation is covered in Partitioning Large Tables for Production Systems.

Can the Partition Key Be Changed Later?

Yes, but changing it on a large production dataset can be expensive.

The process may require creating a new partitioned structure, migrating historical rows, synchronizing writes during migration, validating the result, and switching traffic. The larger the table, the more important it becomes to plan partitioning around stable long-term access patterns.

Conclusion

Database partitioning divides a large logical table into smaller physical units based on a partition key. Used correctly, it improves partition pruning, retention, maintenance, and the operational manageability of very large datasets.

The key decisions are the partition key, partitioning strategy, partition size, indexing model, and lifecycle process. Partitioning works best when those decisions follow actual query patterns and data-retention requirements.

The core principle is: partition data along boundaries that important queries and operational workflows can actually use.

Comments (0)