Table Partitioning

⭐ Interview Importance: MEDIUM
⏱️ Revision Time: 3 min

Concept

As a database table grows past 100 million rows (and into the billions), even B-Tree indexes start to struggle. The index becomes so massively deep and bloated that it cannot fit into RAM. Updating the table becomes incredibly slow.

Table Partitioning is a technique where you take one massive logical table (e.g., orders) and physically chop it into dozens of smaller, manageable tables on the hard drive (e.g., orders_2022, orders_2023, orders_2024).

To the application (Node.js), it still looks and acts like one single orders table. You write SELECT * FROM orders, and the database engine acts as a smart router, silently redirecting the query to the correct physical sub-table.

Partitioning Strategies

1. Range Partitioning (Time-Series)

The most common strategy. You partition the data based on a range of values, almost always dates.

  • logs_january
  • logs_february
  • logs_march

If you query WHERE created_at = '2023-02-15', the database instantly ignores January and March and only searches the February partition. This is called Partition Pruning.

2. List Partitioning

Partitioning based on specific distinct values.

  • users_north_america
  • users_europe
  • users_asia

3. Hash Partitioning

If you don’t have a logical date or region, but you just want to split a massive table into 10 equal chunks to spread the Disk I/O load, you partition based on a mathematical hash of the Primary Key. (e.g., hash(id) % 10).

The Massive Benefit: Data Deletion

In standard SQL, deleting 10 million old log rows is a catastrophic event. You write DELETE FROM logs WHERE created_at < '2020-01-01'. The database executes this row-by-row, locking the table, generating 10 million Dead Tuples (in PostgreSQL), bloating the transaction log, and causing massive Vacuum CPU spikes.

With Range Partitioning, deleting old data is instant and free.
You simply run DROP TABLE logs_2019.
The database deletes the physical file from the hard drive in 10 milliseconds. Zero locks, zero dead tuples, zero performance impact.

Trade-Offs

  • Pros: Massive performance gains for time-series data. Keeps active indexes small and entirely in RAM (e.g., only the index for the “Current Month” partition needs to be cached). Instant data archiving/deletion.
  • Cons: If you write a query that does not include the partition key (e.g., you search for an order by order_id but forget to include the created_at date), the database cannot do Partition Pruning. It is forced to blindly execute the query against every single partition simultaneously, which is drastically slower than searching one massive unpartitioned table.

Interview Questions

Q: Explain the difference between Table Partitioning and Database Sharding.
A:

  • Partitioning splits a table into smaller physical pieces, but all the pieces remain on the exact same physical database server (the same hard drive and the same CPU). It solves localized index bloat and data archiving.
  • Sharding splits a table into smaller physical pieces, and moves those pieces to completely different physical servers (e.g., Server A in New York, Server B in London). Sharding solves hard hardware limits (CPU, RAM, total disk capacity) by distributing the load across a cluster of machines.

Q: You partitioned your sales table by year (2022, 2023, 2024). You run SELECT SUM(amount) FROM sales;. How does the database execute this?
A: Because there is no WHERE clause specifying a date, the database cannot perform Partition Pruning. It will execute the SUM(amount) query across all three physical partitions. In modern databases (like PostgreSQL 11+), it will actually spawn parallel background worker threads to scan all three partitions concurrently, and then sum the three sub-totals together before returning the final result to the user.