You are converting a very large table to a partitioned one and must commit to a partition key and granularity. How do you choose the column, and what makes a partition key a bad one?
answer
- Lifecycle first, predicates second, size third
- Key must be in WHERE, as a literal or parameter
- Immutable key — updates move rows
- Tens to low hundreds of partitions
- Key inside the primary key
basics
~20 sPick the column that appears in most query predicates and governs how data expires — usually a timestamp. Size partitions so each holds a manageable slice, typically tens to a few hundred per table. A bad key is one your queries do not filter on, or one that creates thousands of tiny partitions.
solid answer
~60 sChoose the key from three signals, in priority order: 1. **Data lifecycle.** If rows expire by age, partition on the timestamp — dropping a partition replaces a mass DELETE, which is the biggest single win partitioning offers. 2. **Query predicates.** The key must be in the WHERE clause of your important queries; otherwise every query reads every partition and you have paid overhead for nothing. It must be a column the application actually knows, not one derived by a join. 3. **Even, bounded partition size.** Aim for partitions that are individually manageable — small enough that index maintenance, statistics and rewrites on one are cheap. Bad keys: a column absent from predicates (a status flag when everything is queried by date); a mutable column, since updating it means moving the row between partitions; something so fine-grained that you get thousands of partitions, which inflates planning time, catalog size and per-partition locking; something so coarse that a partition is still hundreds of gigabytes and buys nothing. Granularity follows volume: monthly for tens of millions of rows a month, daily when a day is already large. And the key must be in the primary key and every unique constraint.
go deeper
Say the key should be the column queries filter on, usually a date, and that the key must be in the primary key.
Give the three inputs in order, pick a concrete granularity for a stated volume, and name the mutable-key and too-many-partitions failures.
Argue granularity from real numbers — rows per interval, index size versus memory, maintenance window — and describe how you would validate the choice against the top queries before committing.
Treat it as a long-lived, near-irreversible commitment: reason about how volume and access patterns will change over years, what the migration path is if the key proves wrong, and whether partitioning is warranted at all.
## The decision, framed properly A partition key is close to irreversible: changing it means recreating the table and rewriting every row. It deserves the same care as a primary key. There are exactly three inputs. ## 1. How the data dies Start here, because it is the benefit that no index can replicate. If your retention policy says "keep 13 months", then partitioning by month turns nightly expiry from a `DELETE` that scans, writes tombstones/undo, bloats indexes and generates enormous WAL, into a metadata operation that removes one child table. On a table doing hundreds of millions of rows a month, this alone justifies the project. If the data never expires, this input is silent and you fall through to the next. ## 2. What the queries filter on Partitioning only helps reads when the engine can eliminate partitions, and it can only eliminate them when the query constrains the partition key with a usable predicate. So the key must be a column that is present, as a literal or a parameter, in the predicates of your high-volume queries. Two traps here. First, choosing a column the application does not carry: if a query knows `order_id` but the table is partitioned on `customer_id`, that query reads every partition. Second, choosing a column that queries constrain only through a function or a join — the engine needs a direct comparison on the key to reason about bounds. When queries split between two axes (say, by tenant *and* by time), that is a signal to consider sub-partitioning rather than to compromise on one. ## 3. How big the pieces come out Granularity is a balance between two failure modes. *Too coarse*: yearly partitions on a table doing a billion rows a year means each partition is still enormous, its indexes still too large to cache, and retention still drops a whole year at a time — you have gained almost nothing. *Too fine*: hourly partitions on a five-year retention window is 43,800 child tables. Every query's planning must consider them (even with pruning, the planner and the executor pay per-relation costs), the system catalog grows, each partition needs its own indexes and statistics, locks multiply during DDL, and connections consume more memory holding relation metadata. Most practitioners keep the count in the tens to low hundreds, and treat a few thousand as the point where the overhead starts to hurt. A workable heuristic: pick the interval such that one partition's hot indexes fit comfortably in memory and a maintenance operation on a single partition completes in a time you would be willing to wait — often that lands on monthly or daily for time-series. ## What makes a key actively bad - **Not in the predicates.** The commonest mistake. It converts partitioning into pure overhead: more relations, more planning, more index objects, no elimination. - **Mutable.** If the partition key can be updated, an update that changes it must delete the row from one partition and insert it into another (engines do this for you, but it is a physical row move, it invalidates index entries, and it can fail if no partition accepts the new value). Choose keys that are set at insert and never change — timestamps of the event, tenant of the record. - **Skewed for RANGE or LIST.** If 80% of rows fall in one interval or one listed value, that partition is the table again and the rest are noise. - **Too many distinct values used directly.** Partitioning per customer id with 50,000 customers means 50,000 partitions. Use HASH on that column instead, or partition by a coarser attribute. - **Nullable without a plan.** A NULL key has no home under RANGE/LIST unless a default partition exists, and a default partition that collects everything unmatched can quietly become the biggest one. - **Excluded from the primary key.** Not allowed by the engine for good reason: uniqueness is enforced by per-partition indexes. ## Composite keys and sub-partitioning The partition key can be several columns. Under RANGE they are compared lexicographically, which is useful for `(tenant_id, occurred_at)` when tenants are few and large. More often the cleaner expression is a two-level layout: range by month, then hash or list within each month by tenant. Add the second level only when a real query or maintenance pattern demands it, because it multiplies the partition count. ## Deciding, in practice Write down your top ten queries and your retention rule before choosing. If retention is by time and most queries carry a time window, the answer is range on time at a granularity that yields tens to low hundreds of partitions, and you are done. If retention is unbounded and queries are point lookups by an identifier, hash on that identifier — but ask first whether you need partitioning at all, since a well-indexed table may serve those lookups fine. ## How to answer Give the three inputs in priority order (lifecycle, predicates, size), commit to a concrete key and granularity for a stated table, and name the failure modes: a key absent from predicates, a mutable key, and a partition count in the thousands.
- What actually goes wrong when a table has several thousand partitions?Costs that are per-relation start to dominate: query planning must consider or prune many relations, the system catalog and per-connection relation caches grow, each partition carries its own indexes and statistics that must be maintained, and DDL or maintenance takes locks on many objects at once. Even with good pruning, short OLTP queries can see planning time exceed execution time, which is why practitioners target tens to low hundreds and treat thousands as a design smell.
- Is it a problem if the partition key can be updated?Yes. Changing the key value moves the row physically from one partition to another, which the engine implements as a delete plus insert — invalidating index entries, generating more write traffic, and failing outright if no partition accepts the new value. It also breaks the mental model that a partition is a stable unit you can detach or archive. Prefer keys fixed at insert time, such as the event timestamp or the owning tenant.
saying these in an interview costs you the question
- Choosing a column that queries never filter on and expecting a speedup
- Picking daily or hourly granularity with multi-year retention, producing thousands of partitions
- Partitioning by a mutable column such as order status
- Declaring a unique constraint that does not include the partition key
- Believing partitioning removes the need for indexes inside each partition