Bucketing a Hive table is a promise about where rows live: every row whose join key hashes to bucket 7 is in file 7. When both sides of a join keep compatible promises, Hive can join bucket to bucket and skip redistributing the big table across the cluster. When the promises do not match, Hive quietly runs an ordinary shuffle join and you pay the write-time cost of bucketing for nothing.
The case for bucketing, how to choose the bucket count and the cost model are covered in Hive bucketing. This article covers what happens at query time: how the bucket id is computed, the two join strategies bucketing enables on Tez, the exact conditions and settings each needs, how to read EXPLAIN to prove which one ran, and the failure modes, including one where mixing bucketing versions produced wrong results.
The bucket id: one formula, two hash functions
When Hive writes a bucketed table, each row goes to bucket (hash(key) & Integer.MAX_VALUE) % numBuckets, where the key is the CLUSTERED BY columns. The mask clears the sign bit so negative hashes still give a valid index. Each bucket is written by one task, and its file name carries the bucket number, so a reader knows which keys can be in which file without opening it.
The hash function depends on the table's bucketing_version property. Version 1 is the original Java-style hash, where an INT or BIGINT hashes to roughly its own value. Version 2, introduced in HIVE-18910 for Hive 3, uses Murmur3, which spreads values far more evenly. New tables created by Hive 3 default to version 2, and tables created earlier keep version 1 so their existing files remain valid.
The difference matters in practice. With version 1 and an integer key, bucket assignment is just the key modulo the bucket count, so any pattern in the keys becomes a pattern in the buckets. If every customer ID is a multiple of 4 and the table has 64 buckets, only 16 buckets ever receive rows. Murmur3 removes that coupling, which is the main reason to prefer version 2 for new tables.
Why bucket-aligned data makes joins cheap
A distributed equi-join must bring rows with equal keys to the same task. Without layout knowledge, Hive either broadcasts a small table to every task or shuffles both tables by the join key, which writes, transfers and sorts both inputs. For a fact table of billions of rows, that shuffle dominates the query.
If both tables are bucketed on the join key with the same hash function, rows with equal keys are already in corresponding buckets. With equal counts, bucket b of one table can only match bucket b of the other. With counts where one is a multiple of the other, big-table bucket b can only match small-table bucket b % smallCount. Each task therefore needs only a slice of the other table, and the big table never moves.
Worked example: orders and customers
Take an orders fact table, partitioned by day and bucketed on cust_id into 64 buckets, and a customers dimension bucketed on the same key into 16. Both are ORC, both are sorted on the key within each bucket, and both use bucketing version 2.
-- Both tables bucketed on the join key, with the same hash version.
CREATE TABLE customers (
cust_id BIGINT, name STRING, region STRING
)
CLUSTERED BY (cust_id) SORTED BY (cust_id) INTO 16 BUCKETS
STORED AS ORC
TBLPROPERTIES ('bucketing_version' = '2');
CREATE TABLE orders (
order_id BIGINT, cust_id BIGINT, amount DECIMAL(12,2), order_ts TIMESTAMP
)
PARTITIONED BY (order_date DATE)
CLUSTERED BY (cust_id) SORTED BY (cust_id) INTO 64 BUCKETS
STORED AS ORC
TBLPROPERTIES ('bucketing_version' = '2');
-- Load through INSERT ... SELECT so Hive hashes rows into buckets.
SET hive.exec.dynamic.partition.mode=nonstrict; -- fully dynamic partition insert
INSERT OVERWRITE TABLE orders PARTITION (order_date)
SELECT order_id, cust_id, amount, order_ts, CAST(order_ts AS DATE)
FROM staging_orders;For a query joining one day of orders to customers, each Tez task reads one orders bucket and needs only the customer rows in the matching bucket, one sixteenth of the dimension. If the dimension is 40 GB, a broadcast join would ship all 40 GB to every task and probably exceed the map-join memory budget, while the bucket map join ships about 2.5 GB to each task and builds a hash table only for that slice.
Two details in the DDL matter. The bucket column has the same type on both sides, BIGINT, because the hash is computed on the typed value and an INT and a BIGINT are not guaranteed to hash identically. And rows are loaded with INSERT ... SELECT, which makes Hive hash every row; copying files into the table directory or using LOAD DATA on older versions bypasses that and leaves files that do not honour the layout.
Strategy 1: the bucket map join
A bucket map join is a map join, a hash join with the small side held in memory, restricted to one bucket at a time. On Tez, Hive builds the hash tables from the small table and routes each small-table bucket only to the tasks processing the corresponding big-table buckets, through a custom edge. The big table is read in place.
Its conditions are strict. Both tables must be bucketed on exactly the join columns, with equal counts or counts that are integer multiples. Both must use the same bucketing version. The per-bucket slice of the small side must fit the map-join memory budget, which is what hive.auto.convert.join.noconditionaltask.size controls. The join must be an inner join or an outer join that preserves the big side. When the small table fits in memory as a whole, the planner may prefer an ordinary broadcast map join, which is fine; the bucket map join earns its keep when the small side is too big to broadcast but small per bucket.
-- Map joins in general (on by default in modern Hive)
SET hive.auto.convert.join=true;
SET hive.auto.convert.join.noconditionaltask.size=1000000000; -- bytes, small-side budget
-- Bucket map join on Tez
SET hive.convert.join.bucket.mapjoin.tez=true;
SET hive.optimize.bucketmapjoin=true;
-- Sort-merge-bucket join
SET hive.optimize.bucketmapjoin.sortedmerge=true;
SET hive.auto.convert.sortmerge.join=true;
EXPLAIN
SELECT c.region, SUM(o.amount)
FROM orders o JOIN customers c ON o.cust_id = c.cust_id
WHERE o.order_date = DATE '2026-09-30'
GROUP BY c.region;
Strategy 2: the sort-merge-bucket join
If both tables are also SORTED BY the join key within each bucket, Hive can do better than hashing: it opens corresponding bucket files from both sides and merges them like the merge step of merge sort, advancing whichever side has the smaller key. No hash table is built, so memory use is small and constant regardless of table size, and neither side is shuffled.
A sort-merge-bucket (SMB) join is the right tool when both sides are large, for example two fact tables joined on a shared key, which no hash join can hold in memory. Its conditions add to the bucket map join's: both sides sorted on the join key in the same order, and in practice the same bucket count. Its weakness is skew within a bucket, because one task merges each bucket pair and a hot key makes that task the straggler. Skew handling is covered in skew join optimization.
The planner considers SMB when hive.auto.convert.sortmerge.join is on. If one side is small enough, hive.auto.convert.sortmerge.join.to.mapjoin allows it to switch to a map join instead.
Proving the plan with EXPLAIN
Never assume a bucket join ran. Run EXPLAIN on the real query, with the real partitions, and read two things. The first is the edge types between Tez vertices. A BROADCAST_EDGE sends the whole small table to every task, an ordinary map join. A CUSTOM_EDGE into the map-join vertex is the bucket map join's bucket-aware routing. A SIMPLE_EDGE into a reducer that performs the join means a shuffle. The second is the join operator: a map-join operator inside the big table's map vertex, versus a merge-join operator, and whether the plan mentions bucket or sorted-merge information.
Operator names and the exact wording of plan annotations vary across Hive versions and distributions, so record the plan once when the join works as intended and diff later plans against it. Then confirm at runtime: in the Tez UI, the big-table vertex should not have an outgoing shuffle edge for the join, and the number of tasks should line up with the bucket count of the big table. The Tez execution model behind these vertices and edges is explained in Hive on Tez.
Failure modes
- Mixed bucketing versions. A version-1 table joined to a version-2 table looks bucket-compatible by count but distributes keys differently. HIVE-22098 is titled as data loss when tables with different bucketing versions are joined. Check
DESCRIBE FORMATTEDon both tables and keep versions equal. When a table is rebuilt with CTAS or insert-as-select from an old table, check the version of the result; HIVE-20164 tracked data being hashed with the old logic in that path. - Layout drift. Files copied in by external tools, older
LOAD DATAbehaviour, or writers that do not follow Hive's bucketing leave files that break the promise. A bucket join on such a table can silently miss matches. Spark's native bucketed tables use their own hash and file layout and should not be treated as Hive-bucketed; Iceberg tables use partition transforms instead. - Mismatched counts. 64 and 48 are not multiples, so no bucket join is possible and the query falls back to a shuffle with no error.
- Type mismatches. Joining
INTtoBIGINTorSTRINGtoBIGINTinserts a cast, and the join key is no longer the bucketed value. - Key subsets. A table bucketed on
(cust_id, region)cannot do a bucket join oncust_idalone, because the hash covers both columns. - Skewed buckets. Low-cardinality or patterned keys produce a few huge buckets and long straggler tasks; version 2 fixes patterned integers but not genuinely hot keys.
A small script that compares each table's files with its declared bucket count catches drift early. The file-name pattern varies by version and by ACID layout, so adapt the regular expression to your warehouse:
# Verify that a table's files really match its declared layout before trusting bucket joins.
# Reads bucket ids from ORC file names (bucket_00000, 000000_0, ...) and the declared count.
import re, subprocess, sys
table_dir, declared = sys.argv[1], int(sys.argv[2])
ls = subprocess.run(["hdfs", "dfs", "-ls", "-R", table_dir],
capture_output=True, text=True, check=True).stdout
ids = set()
for line in ls.splitlines():
name = line.rsplit("/", 1)[-1]
m = re.match(r"(?:bucket_)?0*(\d+)(?:_\d+)?(?:_copy_\d+)?$", name)
if m:
ids.add(int(m.group(1)))
print(f"declared={declared} distinct_bucket_files={len(ids)} max_id={max(ids, default=-1)}")
if ids and max(ids) >= declared:
print("bucket id beyond declared count: layout was written with a different count")
Operating bucketed join tables
Treat the bucket layout of a join pair as a contract owned by one team. Record the join key, type, bucket counts, sort order and bucketing version in the table's documentation, and add the layout check to the pipeline that writes the table. Changing any of them requires rewriting the table, so plan migrations as a full rebuild into a new table followed by a swap.
Keep ORC or Parquet file sizes healthy: a 64-bucket table with 365 daily partitions has more than 23,000 files a year, and small-file overhead can cancel the gain, so size the count against the daily data volume, not just the total. Columnar formats also help the merge path because ORC stripes and indexes let readers skip data; see ORC in Hive. If Impala reads the same tables, do not assume its planner exploits the Hive bucket layout the same way; check its query profile for exchange operators before relying on it.
Trade-offs
| Strategy | Needs | Memory | Best for |
|---|---|---|---|
| Broadcast map join | Small side fits in memory | Whole small side per task | Dimensions under the map-join budget |
| Bucket map join | Same key, type, hash version; counts equal or multiples | One small-side bucket per task | Large dimensions joined to facts |
| SMB join | As above plus sorted on key, same counts | Small and constant | Fact-to-fact joins on a shared key |
| Shuffle join | Nothing | Spills as needed | Ad hoc joins, any layout |
Bucketing moves cost from read time to write time. Every write must hash and possibly sort rows into the declared buckets, and the layout locks in one join key. It pays off only for joins that repeat on a stable key at scale.
What to do next
- List the three most expensive recurring joins in your warehouse and the key each one uses.
- Run
DESCRIBE FORMATTEDon both sides of each and record bucket columns, types, counts, sort columns andbucketing_version. - Run
EXPLAINon each join and note whether you see a broadcast, custom or simple edge into the join. - Where versions differ, rebuild the older table with version 2 before relying on any bucket-aware join.
- Add a layout check to the pipeline that writes each bucketed table, and alert on bucket files beyond the declared count.
- For one candidate pair, rebuild both sides with matching layout in a test schema and compare wall time and shuffle bytes against the current plan.