Skip to content
Great engineers know 3 min read · bucketBy 2 practice problems ↓

Bucketing

Pre-shuffle tables into buckets so joins and aggregations on the key skip the shuffle.

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

Comfortable with these? Read on.

TL;DR 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.

PySparkSpark SQL · Write bucketed tables
(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

The plan has Exchange hashpartitioning(customer_id, 200) on both sides, then a Sort, then SortMergeJoin.

Bucketed · no exchange

Bucket 7 of orders holds exactly the keys of bucket 7 of customers. Spark reads them in one task: no Exchange on either side. With sortBy, the Sort can be skipped too when each bucket is one file.

Check it · explain

Run 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.enabled is 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

ConcernDetail
FilesEach 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 countChanging 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 tablesNot 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.
Why this matters: for Delta tables, liquid clustering and Z-ordering co-locate related keys for data skipping, which is the more common need. Bucketing is about avoiding the shuffle in repeated joins. In many lakehouses, broadcast joinsbroadcast join: A join where the small table is copied to every machine, so the big table never has to be shuffled. Learn more → and AQEAQE: Adaptive Query Execution: Spark re-plans a running query using the real data sizes it has measured. Learn more → cover most cases, and bucketing is used only for the few very large, frequently joined tables.

Common mistakes

Bucketing and reading by path

The bucket spec is in the catalog; a path read loses it and shuffles anyway.

Different bucket counts

Spark shuffles one side unless coalescing applies.

Writing without repartitioning

Thousands of tiny files per write.

Bucketing small tables

A broadcast join is simpler and faster.

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 questions

1. 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.

Solve: Pick the Join Strategy →

Keep going

Up next · lesson 20 of 24 · 3 min read
Spark Connect
A thin client that talks to a remote Spark server over gRPC, and what it changes for applications.

Related lessons

Previous: Dynamic partition pruning

Primary sources: Bucketing (DataFrameWriter.bucketBy) · Performance tuning