skip to content

Why partition a lakehouse table by a hash bucket of user_id rather than by user_id itself?

level: middleimportance: should knowfreq 45%

answer

  1. cardinality decides partition count
  2. a fixed number of slots, chosen by you
  3. the same hash is applied to the query literal
  4. order is destroyed, so ranges lose

basics

~10 s

Hash bucketing maps a high-cardinality key onto a fixed number of partitions. Equality lookups on user_id still skip almost everything, while the table keeps a bounded number of partitions holding properly sized files.

solid answer

~40 s

Partitioning directly by a key with tens of millions of distinct values produces tens of millions of partitions, each holding a handful of rows — planning cost and metadata volume then dwarf any scan saving. A bucket transform hashes the key into a fixed count of buckets, say 64 or 1024, so the partition count is a design choice rather than a property of the data. Pruning still works for the access pattern that matters: for `WHERE user_id = 91234` the engine hashes the literal, computes the bucket and reads only that bucket. What you give up is range pruning — hashing deliberately destroys order, so `user_id BETWEEN` prunes nothing. Bucketing is the tool for **high-cardinality equality lookups**; time transforms remain the tool for ranges.

code

text · 9 lines
text
-- 50M distinct users, identity partitioning
users=1/     41 KB
users=2/     18 KB
...  50,000,000 partitions, avg ~100 KB

-- bucket(64, user_id) under day(event_ts)
event_day=2024-05-01/user_bucket=0/   512 MB
event_day=2024-05-01/user_bucket=1/   498 MB
...  64 partitions per day, properly sized files

go deeper

for a junior

Recall that a column with millions of distinct values makes a poor partition column, and that hashing it into a fixed number of groups keeps the partition count under control.

for a middle

Explain how the engine prunes an equality lookup by applying the recorded hash to the query literal, and why the same trick cannot work for a range predicate.

for a senior

Be ready to size the bucket count from target file size and partition volume, and to combine a time transform with a bucket to serve range scans and point lookups from one layout.

for a principal

Own the tradeoff that bucket count is a long-lived commitment: changing it re-maps every key, so history stays under the old mapping and planning spans two layouts until the old data ages out.

## Cardinality is the constraint A partition column's cardinality decides how many partitions you get, and partition count decides how much metadata the planner must process and how large each partition's files can be. This gives a simple rule: the number of distinct values must be small relative to the volume of data, so that every partition holds enough rows to justify at least one properly sized data file. Date satisfies this naturally — a few thousand days for years of history. An identifier does not. Partitioning a 5 TB event table by `user_id` with 50 million distinct users yields 50 million partitions averaging 100 KB. Every query, even one that prunes perfectly, first pays to plan over 50 million partition entries, and every read of a single user still opens a tiny file whose overhead exceeds its payload. ## What a transform is A partition **transform** is a function the table records: the partition value is defined as `f(column)` rather than as a column an ingest job maintains by hand. The family in common use across formats is small: - **time transforms** — year, month, day, hour of a timestamp. Order-preserving, so range predicates translate. - **truncate(width)** — for numbers, round down to a multiple; for strings, keep a prefix. Order-preserving, so ranges and prefix matches translate. - **bucket(N)** — hash the value and take it modulo N. Order-destroying; only equality translates. - **identity** — the value itself, which is what plain Hive-style partitioning does. ## What bucketing buys `bucket(64, user_id)` gives exactly 64 partitions no matter how many users exist. Two properties follow: 1. **Bounded partition count.** The count is set by you, not by the data's cardinality, so it cannot drift as the user base grows. 2. **Equality pruning survives.** For `WHERE user_id = 91234`, the engine applies the *same* recorded hash to the literal, gets a bucket number, and reads only files in that bucket — roughly 1/64th of the table. That is a large win for a needle-in-a-haystack lookup, and it costs nothing at write time beyond the hash. A secondary benefit is join locality: two tables bucketed the same way on the same join key can, in engines that support it, be joined bucket-to-bucket without a full shuffle. ## What bucketing costs - **No range pruning.** `user_id BETWEEN 1000 AND 2000` touches every bucket, by design. If your access pattern is ranges over the key, bucketing is the wrong transform. - **No human-readable layout.** `user_bucket=17` tells an operator nothing about its contents. - **N is effectively permanent in practice.** Changing the bucket count re-maps every value, so existing files no longer satisfy the new rule; formats that support layout evolution will apply the new count only to new data, leaving history under the old bucketing. - **Skew is not solved, only spread.** A single pathological key still lands entirely in one bucket. Hashing fixes *cardinality*, not *hot keys*. ## Choosing N Work backwards from file size. Aim for each bucket, within whatever time partition sits above it, to hold at least one target-sized data file. If a day holds 64 GB and you want files in the hundreds of megabytes, a couple of hundred buckets is defensible; 4096 is not. Powers of two are conventional and make the arithmetic obvious, not because the hash requires it. ## The usual combination Production tables rarely bucket alone. The common shape is a time transform for the range access pattern and a bucket for the point-lookup access pattern: ``` partition by day(event_ts), bucket(64, user_id) ``` A dashboard filtering a date range prunes on the first field; a support tool looking up one user's history in a date range prunes on both. The partition count stays at *days x 64*, which is bounded and predictable. ## Truncate, the middle option Where the key is a string with meaningful prefixes — a country-prefixed code, a sortable identifier — `truncate(width)` keeps order while collapsing cardinality, so both prefix equality and ranges still prune. Reach for it when the key has structure worth preserving; reach for bucketing when it does not.

  • How would you pick the number of buckets for a table receiving 60 GB per day partitioned by day?
    Work backwards from the target file size. At 60 GB per day and files of roughly 512 MB, about 120 files fit in a day, so 64 or 128 buckets keeps one healthy file per bucket per day. Fewer buckets means oversized partitions; many more means the small-file problem returns inside each day.
  • Does bucketing help a join on the bucketed key?
    It can. If both sides are bucketed by the same transform and count on the join key, matching rows are guaranteed to be co-located in the same bucket number, so engines that recognise this can join bucket-to-bucket and skip the shuffle. The requirement is strict: same transform, same count, same key type.
  • One user generates a third of all events. Does bucketing spread that load?
    No. Every row for that user hashes to one bucket, so that bucket is a third of the table. Hashing bounds cardinality, not skew. Fixing a hot key needs a different lever — salting the key, isolating that tenant, or accepting the skew and sizing around it.

Bucketing is a coat check: you do not get one hook per guest, you get sixty-four numbered racks and a rule that always sends the same ticket to the same rack. Finding one coat is fast; finding all coats between two ticket numbers is not.

saying these in an interview costs you the question

  • Says bucketing also lets range filters prune the key
  • Chooses the bucket count without reference to file size
  • Believes hashing removes data skew from a single hot key
  • Thinks the bucket count can be changed freely later
  • Partitions directly by an id column because queries filter on it

context