Two large tables are partitioned the same way on the same key and are frequently joined and aggregated on it. What does joining and aggregating partition-by-partition buy you, what must be true for the optimizer to do it, and why are the settings that enable it (PostgreSQL's enable_partitionwise_join and enable_partitionwise_aggregate) off by default?
answer
- identical bounds + join on the partition key = per-pair joins
- hash table shrinks to 1/N, avoids spilling to disk
- GROUP BY must include the key for the clean aggregate form
- off by default: planning cost scales with partition count
- co-partitioning is a schema commitment, not a knob
basics
~20 sThe optimizer can join matching partition pairs separately and append the results, so each join works on a small slice: hash tables fit in memory, each pair can run in parallel, and aggregation happens per partition. It requires identical bounds and a join on the partition key. It is off by default because planning cost and memory grow sharply with partition count.
solid answer
~60 sWhen both sides are partitioned identically on the join key, a row in one side's January partition can only match the other side's January partition. The planner can therefore split one big join into N small independent joins and append the results. Benefits: each hash table is 1/N the size, so joins that spilled to disk now fit in `work_mem`; the pairs are naturally parallelisable; memory and cache locality improve; and with a matching GROUP BY, aggregation runs per partition and only the partial results are combined, avoiding one enormous hash aggregate. Requirements: the same partitioning strategy, bounds that match (newer PostgreSQL can also handle some non-identical range and list bounds), a join on the full partition key using equality operators from the same operator family, and the aggregate variant additionally needs the partition key among the grouping columns. They are off by default because the planner must consider each partition pair, so planning time and memory grow with partition count. Enable per session or per workload for analytic queries; leave them off for short OLTP queries.
code
text · 9 linesAppend
-> Hash Join
Hash Cond: (o1.customer_id = p1.customer_id)
-> Seq Scan on orders_p0 o1
-> Hash -> Seq Scan on payments_p0 p1
-> Hash Join
Hash Cond: (o2.customer_id = p2.customer_id)
-> Seq Scan on orders_p1 o2
-> Hash -> Seq Scan on payments_p1 p2go deeper
Recognise the concept: if two tables are split the same way on the join column, matching pieces can be joined against each other instead of everything against everything.
State the prerequisites (same strategy, matching bounds, join on the partition key, setting enabled) and the main benefit of smaller hash tables.
Explain the spill-avoidance and parallelism wins, how to verify the plan shape, and why the settings default to off given planning cost.
Judge whether co-partitioning is worth constraining both tables' physical design, weigh planning cost against execution gains at a given partition count, and connect it to node-local joins in a sharded future.
## The idea A join between two partitioned tables is normally planned by appending all partitions on each side and joining the two results. If both tables are partitioned identically on the join key, that is wasteful: the bounds already prove that a row from side A's partition for March can only match side B's partition for March. A partition-wise join exploits that by planning N smaller joins, one per matching pair, and appending their outputs. The same reasoning applies to aggregation. If the GROUP BY includes the partition key, every group lives entirely inside one partition, so the aggregate can run per partition and the results just concatenate. That is a partition-wise aggregate. When the grouping columns do not include the partition key, some engines can still aggregate partially per partition and finalise above the append, which cuts the data flowing upward but requires a combining step. ## Why it matters at scale - **Memory.** A hash join builds a hash table over the smaller input. Split into N pairs, each hash table is roughly 1/N the size. A join that spilled to temporary files at 8 GB may run entirely in memory at 500 MB per pair. Same for hash aggregates and their spill behaviour. - **Parallelism.** Independent per-pair joins map cleanly onto parallel workers, and in distributed or sharded systems the same property is what allows a join to execute locally on each node with no data movement, which is the real reason co-partitioning is a design goal there. - **Locality and pruning composition.** Pruning first reduces the pair list, then each surviving pair is joined against a small, well-cached slice of the indexes. - **Spill avoidance beats raw CPU.** The wins are usually dominated by not spilling to disk rather than by algorithmic savings. ## What must be true 1. **Same strategy, matching bounds.** Both tables range-partitioned by month with the same boundaries, or hash-partitioned with the same modulus. PostgreSQL originally required bounds to match exactly; version 13 added handling for some ranges and lists that do not align perfectly. If one side is partitioned by month and the other by day, or the moduli differ, the optimization does not apply. 2. **The join is on the partition key**, on all of its columns for a composite key, using equality operators from the same operator family as the partitioning. A join on a different column proves nothing about which partitions can match. 3. **The setting is enabled.** `enable_partitionwise_join` and `enable_partitionwise_aggregate` are off by default. 4. **For the aggregate variant**, the partition key must be among the grouping columns for the clean per-partition form. ## Why they are off by default Cost of planning. Without the optimization, the planner considers one join between two appends. With it, it must cost each partition pair, and may consider several join strategies per pair, so planning work and planner memory grow roughly linearly with partition count and can grow worse when several partitioned tables are joined. For a table with hundreds of partitions and a query that runs in a few milliseconds, planning can dominate total latency, and every backend pays that cost. The defaults are therefore chosen for the common OLTP case, and the optimization is opt-in for the analytic case where the join is large enough to repay planning. The pragmatic policy is to enable them narrowly: per session in a reporting connection, per role, or per statement in a batch job, rather than globally in the server configuration. ## The design judgment What makes this a principal-level topic is that it converts a **schema decision** into a **plan capability**. Partition-wise operations are only available if someone deliberately co-partitioned the tables: same key, same strategy, same boundaries, and kept them aligned as the schema evolves. That has real costs, since the partition key must then serve both tables' access patterns, retention policies must line up, and maintenance jobs must keep the boundaries in step. So the question to ask is whether the join in question is on a hot path large enough to justify constraining both tables' physical layout. Sensible reasoning to voice: - If the big join happens once a night in a batch, enabling more parallelism or a larger `work_mem` for that session may be simpler than co-partitioning. - If it happens continuously in an analytic service, co-partitioning plus partition-wise joins can be the difference between spilling and not. - If the data will eventually be sharded across nodes, co-partitioning on the join key is nearly mandatory, because it is what keeps joins node-local. - Watch partition count: the same count that makes per-pair joins small also inflates planning. Coarser partitions often win overall. ## Verifying In the plan, a partition-wise join shows as an Append whose children are join nodes over partition pairs, rather than a single join node above two Appends. A partition-wise aggregate shows aggregate nodes below the Append. Always measure planning time alongside execution time when turning these on, because that is where the regression, if any, will appear.
- You enable enable_partitionwise_join globally and short OLTP queries get slower. Explain and fix.The planner now enumerates and costs each partition pair for every join involving partitioned tables, so planning time and planner memory grow with partition count. For queries whose execution is a few milliseconds, that added planning dominates and every backend pays it. Turn it off globally and enable it narrowly instead, per session for reporting connections, per role, or around the specific batch statements that benefit.
- The two tables are partitioned on the same column but one uses monthly ranges and the other weekly. Can the optimizer still join partition-wise?Generally no. Partition-wise joins need the partition boundaries to correspond so that each partition on one side maps onto a bounded set on the other; PostgreSQL 13 relaxed exact matching for some range and list cases, but months against weeks do not align since a week can straddle a month boundary. The fix is a schema change: align the strategies and boundaries on both tables, which is precisely the co-partitioning commitment this optimization requires.
Instead of pouring two warehouses onto one floor and matching everything, you match aisle 3 against aisle 3: smaller piles, and several teams can work aisles at the same time.
saying these in an interview costs you the question
- Expecting partition-wise joins to work whenever both tables are partitioned, regardless of key, strategy or boundary alignment.
- Turning the settings on globally without measuring planning time on short queries.
- Claiming the join must be on any indexed column rather than on the partition key itself.
- Believing the optimization reduces the total rows joined; it reduces memory, spilling and per-operation size, not the logical work.