Hello! Welcome to your eighth lesson in the "High-Load Distributed Systems" course.
In our last lesson, we established the high-level architectural distinction between single-node partitioning and multi-node sharding. We concluded that partitioning is a powerful technique for managing very large tables on a single server, improving query performance and simplifying data lifecycle management.
Today, we transition from the "why" to the "how." This lesson focuses on the practical design decisions you must make when implementing partitioning. Our goal is to address the learning outcome: Design a partitioning scheme by selecting an appropriate partition key and strategy for a given workload. This is one of the most critical, long-term design decisions you can make for a database, as getting it wrong can be costly to fix later.
We will cover:
- The three native partitioning strategies in PostgreSQL: Range, List, and Hash.
- How to select a partition key that aligns with your query patterns.
- The crucial design principle of a "leading figure" for partitioning interconnected tables in a complex system.
- Key trade-offs, such as the number of partitions versus query planning overhead.
1. The Building Blocks: Partitioning Strategies
PostgreSQL's declarative partitioning offers three distinct strategies for dividing your data. The choice of strategy is dictated by the nature of your data and how you access it.
Documentation: 18: 5.12. Table Partitioning
Let's start with the official PostgreSQL documentation, which provides the canonical definitions for these strategies. Then, we'll look at a more practical article that gives clear examples.
Read the 'Overview' section (5.12.1). Focus on the definitions of Range, List, and Hash partitioning.
When to Consider Postgres Partitioning
Now, let's see these strategies in action with concrete SQL examples from an article by TigerData.
Read the section 'How to Partition a PostgreSQL Table'. This will show you the DDL for creating a parent table and partitions for each of the three strategies.
To summarize the strategies:
-
Range Partitioning: Divides data into partitions based on a continuous range of values. This is the most common strategy, ideal for time-series data (e.g., partitioning by date or timestamp) or sequentially increasing IDs.
This image illustrates range partitioning. A parenttransactionstable is partitioned by year. Queries filtering on a specific year can be routed directly to the corresponding partition, likepublic_transactions_2018, avoiding scans of other years' data. -
List Partitioning: Divides data based on a discrete, explicit list of values. This is suitable for categorical data where the set of possible values is known and finite, such as partitioning a
customerstable bycountry_codeor anorderstable bystatus. -
Hash Partitioning: Distributes data evenly across a fixed number of partitions. PostgreSQL applies a hash function to the partition key and assigns the row to a partition based on the remainder. This is useful when you don't have a natural range or list key but want to spread the write load and data size evenly.
2. The Cornerstone: Selecting a Partition Key
The choice of the partition key is more important than the strategy itself. An effective partition key allows the query planner to perform partition pruning—scanning only the subset of partitions relevant to a query.
The fundamental rule is: Choose a partition key that appears frequently in the WHERE clauses of your most common and performance-critical queries.
Key Constraints and Considerations
There are two critical rules you must follow:
- The partition key must be part of the primary key and any other unique constraints on the table. This is because PostgreSQL enforces uniqueness at the partition level, not globally across all partitions.
- The data type of the column(s) used in the partition key must match the chosen strategy (e.g., a date or numeric type for range, any type for list/hash).
Documentation: 18: 5.12. Table Partitioning
The PostgreSQL documentation provides an excellent list of best practices for choosing a partition key and determining the right number of partitions.
Read section 5.12.6, 'Best Practices for Declarative Partitioning'. Pay close attention to the advice on choosing a partition key and the warnings about having too many partitions.
The Trade-off: Partition Granularity vs. Planning Overhead
As the documentation warns, creating too many partitions can be counterproductive. While smaller partitions can be beneficial for maintenance and data locality, an excessive number increases query planning time. The planner must evaluate each partition's bounds to decide whether to prune it, and this overhead can become the dominant factor in query latency.
Partitioning in Postgres and the risk of high partition counts
This short video from pganalyze provides a compelling case study of this exact problem and its solution.
Watch the entire video. Notice how the initial hash partitioning scheme with many partitions resulted in planning time being longer than execution time. Then, see how ChartMogul solved a similar problem by switching from list partitioning (one partition per customer, leading to thousands of partitions) to a fixed number of hash partitions, resulting in a 5x performance improvement.
A good rule of thumb, as mentioned in the video we'll see next, is to aim for partition sizes in the range of 100-200 GB. This often strikes a good balance between manageability and avoiding excessive planning overhead.
3. Advanced Design: Partitioning an Entire System
In a real-world system, tables don't exist in isolation. You have interconnected tables representing different facets of a business domain. Simply partitioning each table on its "own" best key can lead to disastrous performance for queries that join these tables.
This brings us to the most important concept for system-level partitioning design: the "leading figure."
Partitioning your Postgres tables for 20x better performance | POSETTE 2024
This talk from POSETTE 2024 is perhaps the best explanation of this concept. The speaker uses a real-world example from a financial system to show how misaligned partitions kill performance and how aligning them provides a massive boost.
Watch the segment from 19:04 to 26:00. This is the core of the lesson. Focus on: The 'messy' initial state where each table is partitioned optimally for itself. Why joins across these tables are slow, requiring scans of many partitions. The concept of a 'leading figure' (in their case, transaction_id). The 'clean' final state where all related tables are partitioned by transaction_id, with identical partition boundaries. The dramatic performance improvement because the planner can prune corresponding partitions across all tables in the join.
The key insight is that for joins to be efficient in a partitioned environment, the data being joined must be co-located. By partitioning all related tables (e.g., transactions, payments, transaction_metadata, ledgers) by the same key (transaction_id) and with the same partition boundaries, you ensure that all data for a given transaction range resides in a corresponding set of partitions (e.g., transactions_p1, payments_p1, ledgers_p1).
When you then query for a specific transaction_id, the planner can prune down to a single partition for every table in the join. This transforms a potentially massive, cross-partition join into a small, local join within a single set of co-located partitions.
This design may require some denormalization (e.g., adding the transaction_id to tables that don't naturally have it), a common and necessary trade-off for performance in high-load systems.
4. A Practical Scenario: Designing a Partitioning Scheme
Let's apply these concepts to a hypothetical workload for a high-volume e-commerce platform.
Workload:
- An
orderstable storing billions of records. - An
order_line_itemstable, also with billions of records, linked byorder_id. - Query Pattern 1 (Operational): 80% of queries are from customers checking their recent orders. These queries filter by
customer_idandorder_date(within the last 6 months). - Query Pattern 2 (Analytics): 20% of queries are from the analytics team, scanning all orders for a specific
product_idover the last year. - Data Lifecycle: Order data older than 7 years must be archived and deleted.
Design Process:
-
Identify the Leading Figure: The central entity is the order. However,
customer_idis the most common filter. Let's analyze the trade-offs.- Partitioning by
customer_id(List or Hash): This would be excellent for Query Pattern 1. A query for a specific customer would hit just one partition. However, it would be terrible for Query Pattern 2 (analytics) and for the data lifecycle policy, as a full table scan would be needed for both. - Partitioning by
order_date(Range): This directly supports the data lifecycle policy (drop old monthly/yearly partitions). It also helps both query patterns, as they both have a time component. Queries for recent orders will only scan recent partitions.
- Partitioning by
-
Select the Strategy and Key: Range partitioning by
order_dateis the superior choice. It serves the most critical requirements: efficient data archival and pruning for the most common time-bound queries. -
Determine Granularity: The data must be purged after 7 years. Monthly partitions seem like a reasonable starting point. This would result in
7 * 12 = 84partitions for the active data set, which is a manageable number. We would need a process to create new partitions for future months and drop the oldest ones. -
Apply to the System: Both
ordersandorder_line_itemsshould be partitioned byorder_date. This is not a perfect "leading figure" co-location like thetransaction_idexample, as the join key isorder_id. However, since both tables are partitioned on the same time range, a join filtered by a recent date range will still benefit from pruning on both tables simultaneously. To truly co-locate, we could consider a composite partition key like(order_date, order_id), but that adds complexity. For this scenario, partitioning both byorder_dateis a strong and pragmatic start.
Conclusion
You have now learned the principles of designing a robust and performant partitioning scheme. This is a task that blends deep knowledge of the database's capabilities with a thorough understanding of the application's workload.
Key Takeaways:
- Strategy Follows Workload: Choose Range for continuous data (time), List for categorical data (status, region), and Hash for even distribution when no natural key exists.
- The Key is King: The partition key must align with your most frequent query filters to enable partition pruning.
- Think System-Wide: For interconnected tables, identify a "leading figure" (like a transaction or tenant ID) and partition all related tables by that key to ensure co-location and efficient joins.
- Balance is Everything: Avoid over-partitioning. Too many partitions lead to high query planning overhead, negating the benefits of pruning. Aim for a manageable number of reasonably sized partitions.
Preview of the Next Lesson:
Having covered the design principles, our next step is to put them into practice. In the next lesson, "Implement declarative table partitioning in PostgreSQL with automatic partition routing," we will write the actual DDL to create a partitioned table, define its partitions, and see how PostgreSQL automatically routes data. We will also touch on the maintenance scripts needed to manage this structure over time.