Partitioning & Scaling
How a relational database keeps performing when a table grows past what one heap and one set of indexes can handle, and when the whole workload outgrows a single server. Interviewers use this area to test whether you know the engine-level mechanics — partitioning, sharding, tenancy layout, read/write scaling — rather than just the buzzwords.
part ofRelational database conceptsoverview, primer and where to startread it →on this pageshowhide
explore
- Table partitioning: range, list, hash5 questions
- Partition pruning and partition-aware plans4 questions
- Partition maintenance and retention4 questions
- Sharding a relational database: shard keys5 questions
- Cross-shard queries and resharding6 questions
- Multi-tenancy schemes6 questions
- Vertical vs horizontal scaling of an RDBMS5 questions
- Scaling reads vs scaling writes6 questions
- PostgreSQL DBAroleanchors this topic
- AI & Data Scientistrole
- AI Engineerrole
- Backend Developerrole
- Computer Scienceskill
- Cyber Security Expertrole
- Data Analystrole
- Data Engineerrole
- Forward Deployed Engineerrole
- Full Stack Developerrole
- Java Backend Developerrole
- Kotlin Backend Developerrole
- MLOps Engineerrole
- Machine Learning Engineerrole
- Server-Side Game Developerrole
- Software Architectrole
questions
page 2 of 2You need to add an index to a 4 TB partitioned table with 200 partitions on a live production system, without taking a long lock on the whole table. How would you build it?
basics
~20 sBuild it partition by partition. Create the index concurrently/online on each partition, create the parent index as an empty invalid placeholder, then attach each partition index to it; the parent becomes valid once all are attached. Work is resumable, bounded per partition, and can skip cold partitions.
You put a cache such as Redis or memcached in front of the database to cut read load. What failure modes and consistency issues must the design handle?
basics
~20 sPlan for: stale entries after writes (invalidate after commit, keep TTLs short), stampedes when a hot key expires or the cache restarts (single-flight plus jittered TTLs), the database having to survive a cold cache, and a hit-rate collapse from eviction. A cache is not a source of truth.
A single counter row — say the like count on one viral post — is updated thousands of times per second and overall throughput collapses. What is happening inside the engine, and how would you redesign it?
basics
~20 sUpdates to one row serialise on its exclusive row lock, so throughput is capped at roughly one update per lock-hold time — including commit and network round trips — while everyone else queues, and MVCC version chains and index churn make it worse. Fix by removing the single point: sharded counters, an append-only ledger with rollups, or an in-memory counter flushed periodically.
Read/write splitting can be implemented in application code, in the database driver, or in a connection pooler/proxy. Compare those layers, and describe what commonly breaks when a team turns splitting on.
basics
~20 sApplication-level routing is explicit and safest: the code declares an operation read-only and gets a replica connection. Driver and proxy routing infer intent from statements, which breaks on transactions, session state, SELECT ... FOR UPDATE, and writes hidden inside functions. Whatever the layer, routing must be per transaction, not per statement.
In a sharded relational deployment, something must map each request's shard key to a specific database server. Compare putting that logic in the application/client library, in a dedicated proxy tier, and in a directory or lookup service, and say what each buys you.
basics
~20 sClient-side routing embeds the mapping in the application: fastest, no extra hop, but every service must agree on the topology and redeploy when it changes. A proxy tier centralises routing, pooling and cross-shard merging at the cost of a hop and a component to run. A directory service stores explicit per-key placement, allowing arbitrary moves and pinning, at the cost of a lookup and a critical dependency.
A sharded relational cluster is running hot: one shard is saturated while the others have headroom. How do you determine whether the cause is data skew, request skew, or a single hot key, and what are your options for fixing each?
basics
~20 sMeasure three things per shard: bytes/rows stored, requests per second, and the top keys by request count. Even rows with uneven traffic means request skew; uneven rows means data skew; one key dominating means a hot key. Fix by moving buckets, isolating the heavy tenant, or splitting the key.
An engineer proposes partitioning a 60-million-row table because queries feel slow. How do you decide whether a table is actually big enough or shaped right to be worth partitioning, and what does partitioning cost when it is not warranted?
basics
~20 sPartition when the table is too big to maintain as one unit or when data expires in bulk by a key — not merely because queries are slow. Sixty million rows with a slow query usually needs a better index. Unwarranted partitioning adds planning overhead, more index objects, and constraint restrictions.
A relational database server shows plenty of idle CPU and free memory, yet query latency degrades badly after the application tier grows from 20 instances to 200, each holding its own connection pool. Explain what limits the number of usable database connections and how connection pooling changes the picture.
basics
~20 sEach connection costs memory and a scheduling slot, and concurrency beyond the number of cores and disks only adds context switching, lock contention, and cache thrash. Throughput plateaus then falls. Fix it with a small pool per instance, or a shared external pooler that multiplexes many clients onto few server connections.
Before ordering a larger database machine, how do you determine which resource — CPU, memory, or storage IOPS and latency — is actually limiting the server, and which of those a hardware upgrade will genuinely fix?
basics
~20 sMeasure where time goes, not what looks busy. Sustained high CPU with a hot working set in memory means CPU-bound; heavy read IO with a low cache hit ratio means memory-bound; high commit latency or IO wait means storage-bound. Lock waits and single-threaded hotspots are not fixed by bigger hardware.
When designing a sharded transactional schema, how do you decide between duplicating data so every request stays on one shard versus accepting cross-shard fan-out for the secondary access pattern?
basics
~20 sDecide per access path by rate and latency target. Hot user-facing paths must stay single-shard, so duplicate the data and pay write amplification and a consistency window. Rare back-office paths can fan out. Anything analytical belongs in a derived store, not on the shards.
You own dozens of time-partitioned tables across several services. Design the automation that keeps them healthy: creating future partitions, expiring or archiving old ones. What must the job guarantee, how do you monitor it, and which failure modes do you plan for?
basics
~20 sMake it declarative and data-driven: a config row per table (interval, premake count, retention, archive action). The job must be idempotent, singleton-locked, short-transactioned, and lock-timeout bounded with retries. Monitor runway, partition counts and the last successful run, and prefer a proven implementation such as pg_partman over bespoke scripts.
Write throughput is saturating your primary database while the read replicas sit idle. Walk through how you would decide among batching, queueing, bigger hardware, and splitting the data across multiple primaries.
basics
~20 sMeasure what is actually saturated first — commit/fsync, log volume, index maintenance, lock contention, or connections. Then take the cheap levers in order: batch writes, drop redundant indexes, shorten transactions, move non-critical writes to a durable queue, add IOPS. Splitting across primaries is last because it costs cross-shard queries and transactions forever.
How do you decide that a relational system has reached the end of useful vertical scaling and must move to a scale-out architecture, and what evidence would you bring to that decision?
basics
~20 sScale out when the constraint is structural, not purchasable: the single primary's write ceiling, data volume that breaks maintenance and recovery windows, or a cost curve that has gone nonlinear — and only after tuning, caching, replicas, and partitioning have been exhausted. Bring measured headroom, growth rate, and the date the ceiling arrives.
Some workloads query a huge table along two axes — for example by time window and by tenant. Explain sub-partitioning (composite partitioning), when a second level genuinely pays off, and what it costs.
basics
~20 sSub-partitioning partitions each partition again by a second key — typically RANGE by month at the top and HASH or LIST by tenant underneath. It pays off only when both keys appear in queries or maintenance. The cost is multiplicative: levels multiply into total partition count.
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?
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.
Distributed relational systems such as Citus and Vitess ask you to declare, per table, whether it is distributed on a key or replicated to every shard. How would you make that declaration across a real schema, and why does the choice of distribution column matter so much for joins and transactions?
basics
~20 sDistribute the large, tenant-owned tables on the same key so matching rows land on the same shard; replicate small, shared lookup tables to every shard. Colocated tables can be joined and updated in one local transaction; non-colocated ones force cross-shard joins and distributed commits.
showing 31–46 of 46