PySpark Bucketing: Eliminate Shuffles & Turbocharge Big Data Joins

1. Introduction

At petabyte scale, the single most expensive operation in Apache Spark is the shuffle — the cross-network redistribution of data between stages. A typical join between a 500 GB orders table and a 200 GB customers table can generate hundreds of gigabytes of shuffle traffic, saturating the network, spilling to disk, and stretching job times from minutes to hours.

Bucketing is Spark’s answer to this problem. By pre-organizing data on disk into deterministic buckets at write time, Spark can skip the shuffle entirely at read time — turning an O(n log n) sort-merge join into an O(n) co-located read.

This blog explores bucketing from first principles: what it is, how it works internally, when to use it, and how to implement it correctly in PySpark — with real code, architectural diagrams, and gotchas you’ll only discover in production.

2. What is Pyspark Bucketing?

Bucketing (also called hash partitioning in some systems) is a data organization technique that distributes rows across a fixed number of files — called buckets — based on the hash value of one or more columns.

Each bucket holds all rows whose bucketing column(s) hash to the same value modulo N (the bucket count). This guarantee is the key: when two tables are bucketed on the same column with the same bucket count, rows that could possibly join together are always in the same bucket — so no shuffling is needed.


📌 KEY IDEA
Bucketing trades upfront write cost for zero shuffle on every subsequent read, join, or aggregation. It is a write-once, read-many optimization.

2.1 Bucketing vs Partitioning

These two are often confused. Here are the critical differences:

AspectPartitioningBucketing
Column typeLow-cardinality (date, region)High-cardinality (user_id, order_id)
StorageSeparate directories per valueFixed N files per partition dir
Join optimizationPartition pruning onlyEliminates shuffle entirely
Sort guaranteeNoneSorted within each bucket
Read skew riskHigh (popular dates)Low (uniform hash distribution)
Metastore dependencyNot requiredRequired (Hive Metastore)

3. Bucketing Architecture

Understanding how bucketing works under the hood helps you use it correctly and debug it when it doesn’t behave as expected. The following diagram shows the full data flow:

INPUT DATA Raw DataFrames / Parquet / CSV / JSON — unordered, unpartitioned
▼
BUCKETING WRITE bucketBy(N, col).sortBy(col).saveAsTable() — hashes rows into N bucket files per partition
▼
METASTORE Hive Metastore stores bucket metadata: num_buckets, bucket_cols, sort_cols, format
▼
BUCKET FILES ON DISK part-00000-bucket-0001.parquet … part-XXXXX-bucket-N.parquet — pre-sorted per bucket
▼
BUCKETED READ / JOIN Spark reads matching bucket IDs from both tables — no shuffle, no exchange step in DAG
▼
OUTPUT Final aggregated / joined result — delivered directly to reducers

3.1 The Hash Function

When you call bucketBy(N, ‘customer_id’), Spark applies a deterministic hash function to the bucketing column of every row:

# Conceptual logic (not actual Spark source)
bucket_id = hash(row['customer_id']) % N
# Spark uses Murmur3 hash by default
# The same hash function is used at write time AND read time
# This is what makes co-location possible

The table below illustrates how 6 rows are distributed across 4 buckets:

order_idcustomer_idhash(customer_id)bucket = hash % 4
1001C0010x3A7F…Bucket 0
1002C0020x8B2D…Bucket 1
1003C0010x3A7F…Bucket 0
1004C0030xC51E…Bucket 2
1005C0020x8B2D…Bucket 1
1006C0040x1D9A…Bucket 3
💡 TIPRows with the same customer_id always land in the same bucket — that’s the entire point. When the customers table is also bucketed by customer_id into 4 buckets, matching rows are co-located by design.

3.2 File Layout on Disk

After a bucketed write, Spark produces a structured directory layout. For a table with 2 partitions and 4 buckets:

warehouse/orders/
├── part=2024-01/
│ ├── part-00000-<uuid>_00000.snappy.parquet ← bucket 0
│ ├── part-00000-<uuid>_00001.snappy.parquet ← bucket 1
│ ├── part-00000-<uuid>_00002.snappy.parquet ← bucket 2
│ └── part-00000-<uuid>_00003.snappy.parquet ← bucket 3
└── part=2024-02/
├── part-00000-<uuid>_00000.snappy.parquet
├── part-00000-<uuid>_00001.snappy.parquet
├── part-00000-<uuid>_00002.snappy.parquet

3.3 Metastore Registration

Bucketing metadata is stored in the Hive Metastore. Spark cannot use bucketing if it reads from raw files — the table must be a managed or external table registered in the metastore. The metadata includes:

  • Number of buckets (num_buckets)
  • Bucketing columns (bucket_cols)
  • Sort columns (sort_cols) — if sortBy() was used
  • File format (Parquet, ORC, etc.)

3.4 Join Execution Without Bucketing (SortMergeJoin + Shuffle)

# Without bucketing — Spark must shuffle both tables
Stage 1: Scan orders (500 GB) → hash-partition by customer_id → SHUFFLE
Stage 2: Scan customers (200 GB) → hash-partition by customer_id → SHUFFLE
Stage 3: Sort-merge join on co-partitioned data
# Network traffic: ~700 GB across all nodes
# Duration: ~45 minutes on a 20-node cluster

3.5 Join Execution With Bucketing (Zero Shuffle)

# With bucketing — Spark reads matching bucket IDs directly
Stage 1: Scan orders bucket-0 + customers bucket-0 → local merge join
Stage 2: Scan orders bucket-1 + customers bucket-1 → local merge join
...
Stage N: Scan orders bucket-N + customers bucket-N → local merge join
# Network traffic: ~0 GB (local disk reads only)
# Duration: ~8 minutes on the same cluster

🚀 PERF
In this example, bucketing delivers an 82% reduction in job time and eliminates all shuffle network traffic. Real-world gains vary by cluster size and data distribution.

4. Implementing Bucketing in PySpark

4.1 Basic Bucketed Write

from pyspark.sql import SparkSession
spark = SparkSession.builder \
.appName('BucketingDemo') \
.config('spark.sql.sources.bucketing.enabled', 'true') \
.enableHiveSupport() \
.getOrCreate()
# Load your DataFrame
orders_df = spark.read.parquet('/data/raw/orders')
# Write with bucketing
orders_df.write \
.bucketBy(64, 'customer_id') \
.sortBy('customer_id') \
.mode('overwrite') \
.saveAsTable('orders_bucketed')

4.2 Bucketing Both Sides of a Join

# Write customers table with the same bucket count and column
customers_df.write \
.bucketBy(64, 'customer_id') \
.sortBy('customer_id') \
.mode('overwrite') \
.saveAsTable('customers_bucketed')
# Now read and join — Spark detects compatible bucketing
orders = spark.table('orders_bucketed')
customers = spark.table('customers_bucketed')
result = orders.join(customers, on='customer_id', how='inner')
result.explain() # Should show NO Exchange (shuffle) nodes

4.3 Verifying Shuffle Elimination

After setting up bucketing, always verify the physical plan. A successful bucketed join shows no Exchange nodes:

result.explain(mode='formatted')
# GOOD — bucketing working correctly:
# == Physical Plan ==
# SortMergeJoin [customer_id], [customer_id], Inner
# :- Sort [customer_id ASC NULLS FIRST], false, 0
# : +- FileScan parquet [orders_bucketed] (bucketedScan=true)
# +- Sort [customer_id ASC NULLS FIRST], false, 0
# +- FileScan parquet [customers_bucketed] (bucketedScan=true)
# BAD — shuffle still present (bucketing NOT working):
# Exchange hashpartitioning(customer_id, 64), ENSURE_REQUIREMENTS ← problem!

4.4 Bucketing with Partitioning

Bucketing and partitioning can be combined for multi-dimensional optimization: partition on a date column for time-range pruning, and bucket on the join key to eliminate shuffles within each partition:

orders_df.write \
.partitionBy('order_date') \
.bucketBy(32, 'customer_id') \
.sortBy('customer_id') \
.mode('overwrite') \
.saveAsTable('orders_partitioned_bucketed')

4.5 Bucketed Aggregations

Bucketing also eliminates the shuffle for GROUP BY operations when grouping on the bucketing column:

# With orders bucketed by customer_id, this GROUP BY needs NO shuffle
order_totals = spark.table('orders_bucketed') \
.groupBy('customer_id') \
.agg(
F.count('order_id').alias('order_count'),
F.sum('amount').alias('total_spend')
)
order_totals.explain() # No Exchange node — local aggregation only

5. Choosing the Right Number of Buckets

The bucket count is the most impactful configuration decision. Get it wrong and you either produce tiny files (too many buckets) or oversized files (too few buckets).

5.1 General Formula

# Rule of thumb:
# Each bucket file should be 128 MB – 1 GB after compression
# Formula:
# N = total_uncompressed_data_size / target_bucket_size
# Example:
# Table size: 2 TB uncompressed
# Target bucket size: 512 MB
# N = 2,000,000 MB / 512 MB ≈ 3906 → round to power of 2 → 4096
# Powers of 2 are preferred (they hash more uniformly with Murmur3)
# Common values: 64, 128, 256, 512, 1024, 2048, 4096
Table SizeRecommended BucketsTarget File Size
< 10 GBNo bucketing neededUse standard Parquet
10 GB – 100 GB64 – 128128 MB – 1 GB per bucket
100 GB – 1 TB128 – 512256 MB – 1 GB per bucket
1 TB – 10 TB512 – 2048512 MB – 1 GB per bucket
> 10 TB2048 – 8192512 MB – 1 GB per bucket

⚠️ WARN
All tables participating in the same bucketed join MUST have the same bucket count AND the same bucketing column(s). A mismatch causes Spark to fall back to a full shuffle silently.

6. Key Spark Configurations

Configuration KeyDefaultRecommendedPurpose
spark.sql.sources.bucketing.enabledtruetrueMaster toggle for bucketing
spark.sql.sources.bucketing.autoBucketedScan.enabledtruetrueAuto-detect compatible bucket scans
spark.sql.sources.bucketing.maxBuckets100000100000Upper limit on bucket count
spark.sql.shuffle.partitions200Match bucket countShuffle partition count (for non-bucketed ops)
spark.sql.optimizer.dynamicPartitionPruning.enabledtruetrueWorks with bucketing+partitioning
spark.sql.adaptive.enabledtruetrueAQE can further optimize bucketed plans
spark = SparkSession.builder \
.appName('ProductionBucketing') \
.config('spark.sql.sources.bucketing.enabled', 'true') \
.config('spark.sql.sources.bucketing.autoBucketedScan.enabled', 'true') \
.config('spark.sql.adaptive.enabled', 'true') \
.config('spark.sql.shuffle.partitions', '256') \
.enableHiveSupport() \
.getOrCreate()

7. Limitations & Gotchas

7.1 Requires Hive Metastore

Bucketing only works with Spark managed/external tables registered in Hive Metastore. DataFrames read directly from files (spark.read.parquet(…)) cannot leverage bucketing.

7.2 Bucket Count Must Match on Both Tables

Both tables in a bucketed join must have identical bucket counts and bucket columns. Even a count mismatch of 64 vs 128 disables shuffle elimination.

7.3 Data Skew Can Cause Hot Buckets

If your bucketing column has very skewed distribution (e.g., a few customer_ids account for 90% of rows), those buckets become much larger than others. Consider salting techniques for highly skewed keys.

7.4 Writing Is Slower

Bucketed writes force a full sort of the data before writing. Write times are 2x–5x slower than regular writes. Plan for this in your pipeline scheduling.

7.5 AQE Can Interfere

Adaptive Query Execution (AQE) can sometimes override the bucketed scan optimization. If you see unexpected shuffles with AQE on, try setting spark.sql.adaptive.enabled = false temporarily to diagnose.

⚠️ WARNAs of Spark 3.x, bucketing does not work with Delta Lake’s native writer (delta format). Use Parquet or ORC for bucketed tables, or use Delta’s Z-Ordering as an alternative optimization.

8. Production Best Practices

  1. Always verify shuffle elimination with explain() before promoting to production.
  2. Standardize bucket counts across your data platform — pick a small set (64, 256, 1024) and enforce them.
  3. Combine bucketBy with sortBy on the same column — enables merge-sort without re-sort at join time.
  4. Monitor bucket file sizes after writes — rewrite if files drift outside the 128 MB – 1 GB range.
  5. Use partitionBy + bucketBy together for time-series data with frequent joins.
  6. Document bucketing decisions in your data catalog — future readers need to know why bucket counts match.
  7. Schedule bucketed table rewrites during off-peak hours — they are write-intensive operations.
  8. Run ANALYZE TABLE after bucketed writes to update statistics and help the query optimizer.
# Always run statistics collection after bucketed writes
spark.sql('ANALYZE TABLE orders_bucketed COMPUTE STATISTICS')
spark.sql('ANALYZE TABLE orders_bucketed COMPUTE STATISTICS FOR COLUMNS customer_id, amount')

9. When to Use Bucketing

ScenarioUse Bucketing?Reason
Large table joins (>10 GB each side)YESShuffle elimination gives 3x–10x speedup
Frequent GROUP BY on same columnYESAggregation also benefits from co-location
Small table joins (<1 GB one side)NOBroadcast join is faster and simpler
One-time ETL jobNOWrite overhead not worth it for single use
Streaming sink tableMAYBEWorks but micro-batch rewrites are expensive
Highly skewed join keysCAREFULLYUse salting or partial bucketing
Delta Lake tablesNOUse Z-Ordering instead

10. Summary

Bucketing is one of the most powerful — and most underused — performance optimizations in PySpark. When applied correctly, it transforms expensive shuffle-heavy joins into fast, co-located reads that scale linearly with cluster size.

The core principles to remember:

  • Bucketing pre-organizes data at write time so reads require no shuffle.
  • Both tables in a join must have matching bucket counts and columns.
  • Bucketing requires Hive Metastore — it does not work on raw file reads.
  • Always verify with explain() — Spark can silently fall back to shuffling.
  • Combine with partitioning and sortBy for maximum query performance.

With careful planning — choosing the right bucket count, standardizing across your data platform, and monitoring file sizes — bucketing can be the single biggest performance win in your Spark data pipeline.


Discover more from DataSangyan

Subscribe to get the latest posts sent to your email.

Leave a Reply