Partitioning a Table

We need to make a distinction between partitioning and sharding. In a way, sharding means partitioning the data across multiple database servers. However, the term “partitioning” can also refer to segmenting a single large table within a database instance. It is more common to do this in a relational database, but it can also be done in other types of databases.

Keep in mind that some resources use the terms interchangeably. For example, the first edition of Designing Data-Intensive Applications uses “partition” for what we call a shard. The second edition switched to “shard”. DynamoDB’s documentation also uses “partition” for what we call a shard. In Cassandra, a “partition” is the set of rows that share a partition key, which is closer to one shop’s data than to a shard.

In this course, we use partitioning to mean dividing one table into smaller tables, called partitions, inside one database server.

Partitioning

Let’s go back to the Shop API where we had a relational database. Suppose the orders table keeps every order ever placed, so it grows every day. Most queries ask for recent orders, such as yesterday’s. We can use partitioning to manage this growth.

For example, we could divide orders by month:

graph TB
    subgraph S[One Database Server]
        T[orders] --> P1[orders_2026_08]
        T --> P2[orders_2026_09]
        T --> P3[orders_2026_10]
    end

A partition key is the column or combination of columns used to decide which partition holds a row. Here, we use placed_at and assign each order to the partition for its month. An order placed in September 2026 goes into orders_2026_09.

Here is how we could create this table in PostgreSQL, showing only the columns needed for this example:

CREATE TABLE orders (
    id BIGINT NOT NULL,
    placed_at TIMESTAMP NOT NULL
) PARTITION BY RANGE (placed_at);

CREATE TABLE orders_2026_08 PARTITION OF orders
    FOR VALUES FROM ('2026-08-01') TO ('2026-09-01');

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

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

PARTITION BY RANGE (placed_at) tells the database to divide the table by ranges of placed_at. Each PARTITION OF orders statement creates one partition. The FROM boundary is included, and the TO boundary is excluded. We create another partition for each new month before inserting orders for that month.

The application still queries orders. For example, this query reads the orders placed on October 5:

SELECT id, placed_at
FROM orders
WHERE placed_at >= '2026-10-05'
  AND placed_at < '2026-10-06';

We would write the same query if orders were not partitioned. The query does not name orders_2026_10; the database determines which partition to read from the date conditions.

Our example uses range partitioning. Other types of partitioning exist, including list and hash partitioning.

Reading and removing orders

Suppose the application sends a query for orders placed yesterday. The query’s start and end times are both in October 2026. The database only needs to read the October partition. It can skip August, September, and the earlier partitions. Skipping partitions that cannot hold rows matching the query is called partition pruning. Each partition also has its own smaller indexes.

A query for orders from a particular shopper, with no date restriction, may need to read every partition. The database cannot use the month boundaries to rule any of them out.

Partitioning also helps when we remove old data in bulk. Suppose we remove orders older than seven years every month. Once every order in a partition is older than seven years, we can drop that partition. Otherwise, the database would have to delete millions of rows one at a time and then clean up after all of those deletions. Dropping or detaching a partition can be much faster than that.

Partitioning and sharding

Partitioning can reduce the work needed to query and maintain a large table. Sharding distributes data and its workload across servers to increase storage, read, and write capacity.

If queries for recent orders slow down as years of orders accumulate, we should first check their indexes. Partitioning by month can then help the database read only the partitions for the months those queries ask for.