A lakehouse table partitioned by hour and customer_id has millions of partitions — what breaks?
answer
- cost moved from reading to planning
- two cardinalities multiplied together
- perfect pruning and still slow
- choose the partition count instead of inheriting it
basics
~20 sCost moves from scanning to planning: the engine must enumerate and filter millions of partition entries before reading anything, files fall far below target size, and skew grows. Coarsen the time grain and hash the customer key into buckets.
solid answer
~50 sOver-partitioning trades one bottleneck for a worse one. Hour crossed with customer identity multiplies cardinalities, so partition count grows with both time and the customer base, and each partition ends up holding kilobytes rather than hundreds of megabytes. Query latency is then dominated by **planning** — enumerating partitions and file entries, evaluating predicates over them — which is work every query pays even when it prunes perfectly. Reads that survive pruning open many tiny files, so per-file overhead outweighs payload, and partitions are wildly uneven because customers are. The fix is to reduce cardinality on both axes: partition by day rather than hour, and replace the raw customer key with a bounded number of hash buckets, so partition count becomes a chosen constant per day instead of a function of how many customers signed up.
code
text · 9 linesbefore: hour x customer_id
partitions 4,180,000
avg bytes/partition 2.6 MB
planning / execution 31 s / 4 s
after: day x bucket(64, customer_id)
partitions 23,360
avg bytes/partition 465 MB
planning / execution 0.4 s / 6 sgo deeper
Recall that too many partitions is a real failure mode: each one ends up tiny, and the engine spends its time working through a huge list rather than reading data.
Explain that partition count is the product of the partition fields' cardinalities, and that planning happens before any I/O, so a perfectly pruned query still pays it.
An interviewer wants the diagnosis from evidence — partition count, bytes per partition, planning versus execution time — followed by coarsening the time grain and bucketing the high-cardinality key.
Own the policy angle: define what a healthy partition looks like for your platform, decide whether dominant tenants get their own tables, and plan how a layout change rolls out given that it only governs new writes.
## Diagnosing it, not guessing Before prescribing anything, get three numbers: total partition count, average bytes per partition, and — from a representative query's plan — how long planning took versus execution, plus how many partitions survived pruning. Over-partitioning has a distinctive signature: pruning works *beautifully* (a handful of partitions survive) and the query is still slow, because the time went into deciding which handful. If planning is a small fraction of runtime and the surviving partitions are large, you have a different problem and should stop here. ## Why hour x customer_id explodes Partition count is the product of the cardinalities of the partition fields. A year of hours is about 8,760 values; ten thousand customers turn that into tens of millions of possible partitions, and even a sparse table materialises millions of them. Two independent growth axes is the structural mistake: the count grows as history accumulates *and* as the business acquires customers. ## What actually degrades **Planning cost.** Every query begins by narrowing the partition and file list. Whether that list lives in a metastore or in the table's own metadata, filtering millions of entries takes real time and memory, and it happens before a single byte of data is read. This cost is paid per query, including by queries that end up reading almost nothing — which is why dashboards with tight filters feel inexplicably slow. **Tiny files.** A partition receiving one customer's events for one hour may hold a few kilobytes. Reading it costs a request, a footer read, and decoder setup, all to return a handful of rows. The per-file constant dominates, and the effect compounds because a query spanning a day and many customers opens thousands of such files. (The remedy for file size itself — compaction toward a target size — is a separate lever; here the point is that the *layout* guarantees the problem returns after every compaction, because the partition boundaries force the split.) **Skew.** Customer volume follows a power law. Some partitions hold gigabytes and others a few rows, so parallelism is uneven: most tasks finish instantly and a few run long, leaving the cluster idle at the tail. **Operational drag.** Metadata grows, listing and maintenance operations take longer, and anything that must walk the partition list — cleanup jobs, catalog syncs, statistics collection — slows in proportion. ## The opposite failure, for contrast Under-partitioning is the mirror image: too few partitions, each enormous, so a filter that should touch a day instead reads a month. The symptom is the reverse — planning is instant, the scan is huge, and pruning shows most partitions surviving. Knowing both signatures is what lets you say which one you have from a plan rather than from intuition. Between them is a defensible target: partitions large enough to hold at least one properly sized data file, and few enough that the planner's work is negligible next to the scan. ## The fix, in order 1. **Coarsen the time grain.** Hour to day divides the count by 24 immediately. Hour is justified only when a large fraction of queries genuinely bound themselves to a few hours *and* an hour holds enough data to be worth isolating. 2. **Bound the second axis.** Replace `customer_id` identity with a hash bucket of fixed count. Point lookups on one customer still prune to one bucket, but the partition count per day is a constant you chose rather than the size of the customer table. 3. **Or drop the second axis entirely.** If most queries filter time and only some filter customer, partition by day and let finer pruning come from ordering data within the partition by customer, so file-level statistics narrow the scan. That keeps the partition count minimal. 4. **Apply the new layout going forward.** In a format that supports changing the layout in place, this is a metadata commit affecting new writes only; history keeps its old, bad layout until you decide whether it is worth rewriting or simply let retention expire it. ## Isolating the truly large tenant One more option is worth naming for a lead-level answer: if a handful of customers dominate volume, uniform partitioning serves neither them nor the long tail. Splitting those tenants into their own tables — or accepting deliberate skew and sizing the maintenance job around it — is often better engineering than a single layout that is a compromise for everyone. ## What to say in the interview Name the mechanism (planning cost paid per query before any I/O), name the two contributing cardinalities, propose coarsening plus bucketing, and be explicit that the layout change helps new data first. Candidates who answer only "run compaction" have addressed the symptom while leaving the layout that regenerates it in place.
- How do you tell over-partitioning from under-partitioning using only a query plan?Compare planning time with execution time and look at what pruning returned. Over-partitioning shows long planning, a tiny surviving set, and small bytes scanned. Under-partitioning shows instant planning and an enormous scan because the surviving partitions are far broader than the filter. The two prescriptions are opposite, so the distinction matters.
- If most queries filter on time and only some on customer, what layout would you choose?Partition by day alone and keep the partition count minimal, then order rows within each partition by customer so per-file statistics narrow the customer-filtered scans. That serves the common access pattern with no partition-count cost and still gives the occasional customer query useful skipping without a second partition axis.
- Why does running compaction not durably fix this table?Compaction can only merge files within a partition boundary. If the layout guarantees that one customer-hour is its own partition, the merged result is still a small file, and the next hour of writes recreates the pattern. The layout is what regenerates the problem, so the layout is what must change.
- After coarsening the layout, when does query performance actually improve?For new data immediately, and for the table as a whole only as new data accumulates, because existing files keep the layout they were written under. Queries bounded to recent time see the benefit at once; queries over long history keep paying the old planning cost until that data is rewritten or expires.
It is a warehouse with one labelled shelf per parcel: finding a parcel is trivial, but walking the index of ten million shelf labels before you start takes longer than the walk to the shelf.
saying these in an interview costs you the question
- Says compaction alone fixes it, leaving the layout unchanged
- Assumes more partitions always improve pruning
- Ignores that planning cost is paid before any data is read
- Proposes adding a third partition column to narrow further
- Expects a layout change to improve historical data immediately