Bucketing
Pre-shuffle tables into buckets so joins and aggregations on the key skip the shuffle.
On this page
Show code in
Every code block on the page follows this.
You will learn
- What bucketing does to a table's files
- How it removes the shuffle from joins and aggregations
- The conditions both tables must meet
- Why table formats moved to clustering instead
Read first
- Narrow vs wide transformations · 3 min read
- Broadcast joins · 3 min read
- Partitioning done right · 4 min read
Comfortable with these? Read on.
bucketBy(64, "customer_id") hashes rows into 64 buckets at write time and records that in the catalog. A later join or groupBy on customer_id between tables bucketed the same way can skip the shuffle, because matching keys are already in the same bucket. It pays off for large tables joined repeatedly on the same key; it is a Hive-style table feature and needs saveAsTable.What it does
A sort-merge join shufflesshuffle: Moving rows between machines so that all rows with the same key end up together. Needed by joins, groupBy and sorting, and usually the most expensive step of a job. Learn more → both sides so that rows with the same key meet in the same tasktask: The work for one partition in one stage, run on one CPU core. Learn more →. If both tables are joined on customer_id every day, that shuffle is paid every day. Bucketing pays it once, at write time: rows are hashed by the key into N files per partitionpartition: A chunk of a DataFrame's rows. Spark processes each partition as one task, so partitions decide how much work runs in parallel. Learn more →, and the table metadata says so.
(orders.write .bucketBy(64, "customer_id") .sortBy("customer_id") .mode("overwrite") .saveAsTable("sales.orders_b")) (customers.write .bucketBy(64, "customer_id") .sortBy("customer_id") .mode("overwrite") .saveAsTable("sales.customers_b")) joined = spark.table("sales.orders_b").join(spark.table("sales.customers_b"), "customer_id")
CREATE TABLE sales.orders_b USING parquet CLUSTERED BY (customer_id) SORTED BY (customer_id) INTO 64 BUCKETS AS SELECT * FROM orders;
What changes in the plan
Unbucketed · shuffle both
Exchange hashpartitioning(customer_id, 200) on both sides, then a Sort, then SortMergeJoin.Bucketed · no exchange
sortBy, the Sort can be skipped too when each bucket is one file.Check it · explain
joined.explain() and confirm there is no Exchange under the join. If there is, a condition below is not met.Conditions
- Both tables bucketed on the join key (the same columns, same order).
- The same number of buckets, or (Spark 3.1+, with
spark.sql.bucketing.coalesceBucketsInJoin.enabled) one a multiple of the other. - Tables read through the catalog (
spark.table), not by file path: the bucket spec lives in the metastore. spark.sql.sources.bucketing.enabledis true (the default).- Spark's hash function is used for both: tables bucketed by Hive or another engine are not compatible.
Costs and trade-offs
| Concern | Detail |
|---|---|
| Files | Each write task can write one file per bucket: 200 tasks × 64 buckets = 12,800 files. Repartition by the bucket key with the bucket count before writing. |
| Fixed bucket count | Changing it means rewriting the table. Size buckets for the table's future, often 100 MB-1 GB per bucket. |
| Skewskew: When one key has far more rows than the others, so one task does most of the work while the rest wait. Learn more → | A hot key fills one bucket. |
| Only Spark/Hive tables | Not supported by Delta Lake. Iceberg has a bucket(N, col) partition transform, used for layout and pruning, and storage-partitioned joins (Spark 3.3+) can use it to avoid shuffles. |
Common mistakes
Bucketing and reading by path
Different bucket counts
Writing without repartitioning
Bucketing small tables
Key takeaways
- Bucketing pre-shuffles a table by key at write time.
- Joins and aggregations on the bucket key can skip the Exchange.
- Both sides need the same key and compatible bucket counts, read via the catalog.
- Delta uses clustering instead; Iceberg has bucket transforms.
Check yourself
3 questions1. What does bucketing remove from a join on the bucket key?
Show the answer
The shuffle (Exchange). Matching keys are already in matching buckets.
2. Why might a join of two bucketed tables still shuffle?
Show the answer
They were read by file path, so the bucket spec is unknown. Bucketing metadata lives in the catalog.
3. 200 tasks write into 64 buckets without repartitioning. Up to how many files?
Show the answer
12,800. Each task can write one file per bucket.
Practice it
Interview problems that use this: write the PySpark, run it, and get graded on hidden tests.
Keep going
Up next · lesson 20 of 24 · 3 min readSpark Connect
A thin client that talks to a remote Spark server over gRPC, and what it changes for applications.
Related lessons
Broadcast joinsData lake & lakehouse · 3 min read
Liquid clusteringSpark internals · 3 min read
Narrow vs wide transformations
Previous: Dynamic partition pruning
Primary sources: Bucketing (DataFrameWriter.bucketBy) · Performance tuning