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:
| Aspect | Partitioning | Bucketing |
| Column type | Low-cardinality (date, region) | High-cardinality (user_id, order_id) |
| Storage | Separate directories per value | Fixed N files per partition dir |
| Join optimization | Partition pruning only | Eliminates shuffle entirely |
| Sort guarantee | None | Sorted within each bucket |
| Read skew risk | High (popular dates) | Low (uniform hash distribution) |
| Metastore dependency | Not required | Required (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_id | customer_id | hash(customer_id) | bucket = hash % 4 |
| 1001 | C001 | 0x3A7F… | Bucket 0 |
| 1002 | C002 | 0x8B2D… | Bucket 1 |
| 1003 | C001 | 0x3A7F… | Bucket 0 |
| 1004 | C003 | 0xC51E… | Bucket 2 |
| 1005 | C002 | 0x8B2D… | Bucket 1 |
| 1006 | C004 | 0x1D9A… | Bucket 3 |
| 💡 TIP | Rows 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 tablesStage 1: Scan orders (500 GB) → hash-partition by customer_id → SHUFFLEStage 2: Scan customers (200 GB) → hash-partition by customer_id → SHUFFLEStage 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 directlyStage 1: Scan orders bucket-0 + customers bucket-0 → local merge joinStage 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 SparkSessionspark = SparkSession.builder \ .appName('BucketingDemo') \ .config('spark.sql.sources.bucketing.enabled', 'true') \ .enableHiveSupport() \ .getOrCreate()# Load your DataFrameorders_df = spark.read.parquet('/data/raw/orders')# Write with bucketingorders_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 columncustomers_df.write \ .bucketBy(64, 'customer_id') \ .sortBy('customer_id') \ .mode('overwrite') \ .saveAsTable('customers_bucketed')# Now read and join — Spark detects compatible bucketingorders = 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 shuffleorder_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 Size | Recommended Buckets | Target File Size |
| < 10 GB | No bucketing needed | Use standard Parquet |
| 10 GB – 100 GB | 64 – 128 | 128 MB – 1 GB per bucket |
| 100 GB – 1 TB | 128 – 512 | 256 MB – 1 GB per bucket |
| 1 TB – 10 TB | 512 – 2048 | 512 MB – 1 GB per bucket |
| > 10 TB | 2048 – 8192 | 512 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 Key | Default | Recommended | Purpose |
| spark.sql.sources.bucketing.enabled | true | true | Master toggle for bucketing |
| spark.sql.sources.bucketing.autoBucketedScan.enabled | true | true | Auto-detect compatible bucket scans |
| spark.sql.sources.bucketing.maxBuckets | 100000 | 100000 | Upper limit on bucket count |
| spark.sql.shuffle.partitions | 200 | Match bucket count | Shuffle partition count (for non-bucketed ops) |
| spark.sql.optimizer.dynamicPartitionPruning.enabled | true | true | Works with bucketing+partitioning |
| spark.sql.adaptive.enabled | true | true | AQE 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.
| ⚠️ WARN | As 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
- Always verify shuffle elimination with explain() before promoting to production.
- Standardize bucket counts across your data platform — pick a small set (64, 256, 1024) and enforce them.
- Combine bucketBy with sortBy on the same column — enables merge-sort without re-sort at join time.
- Monitor bucket file sizes after writes — rewrite if files drift outside the 128 MB – 1 GB range.
- Use partitionBy + bucketBy together for time-series data with frequent joins.
- Document bucketing decisions in your data catalog — future readers need to know why bucket counts match.
- Schedule bucketed table rewrites during off-peak hours — they are write-intensive operations.
- Run ANALYZE TABLE after bucketed writes to update statistics and help the query optimizer.
# Always run statistics collection after bucketed writesspark.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
| Scenario | Use Bucketing? | Reason |
| Large table joins (>10 GB each side) | YES | Shuffle elimination gives 3x–10x speedup |
| Frequent GROUP BY on same column | YES | Aggregation also benefits from co-location |
| Small table joins (<1 GB one side) | NO | Broadcast join is faster and simpler |
| One-time ETL job | NO | Write overhead not worth it for single use |
| Streaming sink table | MAYBE | Works but micro-batch rewrites are expensive |
| Highly skewed join keys | CAREFULLY | Use salting or partial bucketing |
| Delta Lake tables | NO | Use 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.