Hive can execute the same join in five physical ways, and it picks one per join, at compile time, from estimates. The choice decides whether a query broadcasts a few megabytes or sorts terabytes, and whether it finishes or dies with an out-of-memory error. The algorithms themselves are covered elsewhere on this site: join semantics, shuffle joins, map joins and which side can be broadcast in Hive Joins, in depth, bucket map and sort-merge-bucket joins in Hive bucketing for joins, and join ordering and statistics in the Hive cost-based optimizer.
This page is about the selection itself: what the Tez compiler runs after the join order is fixed, the order of attempts, the numbers compared, and why a small change in statistics or settings moves a query between strategies. It is read from the ConvertJoinMapJoin class and HiveConf on Apache Hive's master branch in October 2026. Defaults quoted are master's; vendor distributions often ship different values, so always read your own session.
Five physical strategies
| Strategy | Data movement | Memory need | Plan marker |
|---|---|---|---|
| Map join | Small sides broadcast to every task of the big side; big side not shuffled | All small sides together, per task | Map Join Operator, BROADCAST_EDGE |
| Bucket map join | Bucket i of the small side goes only to tasks reading bucket i of the big side | One bucket of each small side | Map Join Operator, CUSTOM_EDGE |
| Dynamically partitioned hash join | Both sides shuffled by key without sorting; reducers build hash tables | Small sides divided by reducer count | DynamicPartitionHashJoin: true, CUSTOM_SIMPLE_EDGE |
| Sort-merge-bucket join | None: matching sorted buckets merged in place | Streaming | Merge Join Operator in a Map vertex |
| Merge (shuffle) join | Both sides shuffled and sorted by key | Streaming, plus sort buffers | Merge Join Operator, SIMPLE_EDGE |
The dynamically partitioned hash join, DPHJ, is the least known. Its design discussion observed that most CPU in shuffle joins went into sorting and merging. Tez can shuffle without sorting, so DPHJ partitions both sides by key and hash-joins each partition, provided each reducer's share of the small side fits in memory.
The order of attempts
The selection runs once per join operator. In outline it tries the cheapest strategies first and falls back step by step. The pseudocode below follows the control flow of the process method; helper names are simplified.
choose_physical_join(join): # one call per join operator, Tez
if not hive.auto.convert.join:
return SMB if smb_possible(join) else MERGE_JOIN # DPHJ also needs auto.convert.join
limit = maxJoinMemory # noconditionaltask.size, inflated on LLAP
buckets = estimate_buckets(join) if hive.convert.join.bucket.mapjoin.tez else 1
mj = map_join_conversion(join, divisor=buckets, limit=limit, check_thresholds=True)
if mj is None:
return SMB if smb_possible(join) else reduce_side(join)
if buckets > 1:
if llap:
if not dphj_enabled:
if bucket_map_join_ok(join, mj): return BUCKET_MAP_JOIN
elif small_side_total(mj) > limit: # otherwise a plain map join is fine
if bytes(small + big) < nodes * bytes(small):
if dphj_ok(join): return DPHJ
elif bucket_map_join_ok(join, mj): return BUCKET_MAP_JOIN
elif bucket_map_join_ok(join, mj):
return BUCKET_MAP_JOIN
mj = map_join_conversion(join, divisor=1, limit=limit, check_thresholds=True)
if mj is None or (mj.full_outer and not full_outer_map_join_supported):
return reduce_side(join)
return MAP_JOIN
reduce_side(join):
if hive.auto.convert.join and hive.optimize.dynamic.partition.hashjoin and dphj_ok(join):
return DPHJ # divisor = reducers, entry/shuffle checks off
return MERGE_JOINThree details surprise people. SMB is tried only when no map join is possible, even for perfectly bucketed and sorted tables. With hive.auto.convert.join off, DPHJ is off too. And the bucket count is a divisor, so a small side too big to broadcast can still qualify when only one bucket at a time must fit.
The numbers being compared
Every comparison uses an online data size estimate: the operator's estimated data size after filters and projection, adjusted for how the hash table stores values, plus a per-row overhead and a per-slot overhead. The slot count is the row estimate divided by hive.hashtable.loadfactor, default 0.75, rounded up to a power of two, so row count matters as much as bytes.
The budget is maxJoinMemory. Outside LLAP it is simply hive.auto.convert.join.noconditionaltask.size, whose master default is 10,000,000 bytes. On LLAP it is inflated because a query may borrow memory from other executors: the budget becomes size + size × factor × slots, where the factor is hive.llap.mapjoin.memory.oversubscribe.factor, default 0.2, and slots is bounded by executors per node and hive.llap.memory.oversubscription.max.executors.per.query. With a 10 MB budget and three slots, the effective limit is 16 MB.
How the big table is chosen
A map join streams one input, the big table, and hashes the others. The join type limits which positions may be streamed; in a left outer join only the left side can be. The loop walks the inputs in order:
def map_join_conversion(inputs, candidates, divisor, limit):
"""Python transliteration of the big-table loop in getMapJoinConversion.
inputs: list of {"bytes": in-memory estimate, "card": cumulative row estimate}
candidates: positions the join type allows to be streamed (the big table)
Returns the big-table position, or None if no map join is possible.
(The hashtable.max.entries / shuffle.max.size check is left out here.)
"""
big, big_bytes, big_card = None, -1, -1
overflow, total = False, 0
for pos, inp in enumerate(inputs):
size, cur_overflow = inp["bytes"], False
if big is None or size > big_bytes:
if overflow:
return None # a second input that cannot fit in memory
if size / divisor > limit:
if pos not in candidates:
return None # too big, and the join type forbids streaming it
cur_overflow = overflow = True
card = -1 if overflow else inp["card"]
chosen = pos in candidates and (
big is None or cur_overflow or
(not overflow and (card > big_card or (card == big_card and size > big_bytes))))
if big is not None and chosen:
total += big_bytes # the previous big table becomes a hash table
elif not chosen:
total += size
if total / divisor > limit:
return None # small sides together exceed the budget
if chosen:
big, big_bytes, big_card = pos, size, card
return bigAn input that cannot fit in memory, after dividing by the bucket count, must become the big table; a second such input ends the attempt. Otherwise the big table is the candidate with the highest cumulative row estimate, size breaking ties. Everything else counts toward the hash-table total, which must stay within budget. Missing statistics on any input mean no conversion at all, which is why tables without statistics so often end up in a merge join.
One more check sits after the loop. If any hashed input is predicted to hold more entries than hive.auto.convert.join.hashtable.max.entries, master default 21,000,000, and the big side is no larger than hive.auto.convert.join.shuffle.max.size, master default 10,000,000,000 bytes, the map join is refused and the join goes to DPHJ if enabled, otherwise to a merge join. A big side above that shuffle limit keeps the map join, since shuffling it would cost more than a crowded hash table.
Bucket map join versus DPHJ on LLAP
When the inputs carry bucketing and hive.convert.join.bucket.mapjoin.tez is on, the first attempt divides sizes by the bucket count. Outside LLAP, a successful attempt becomes a bucket map join if the bucket layout allows it. On LLAP, if DPHJ is disabled the compiler tries the bucket map join directly. If DPHJ is enabled, it first checks whether the small sides fit the budget without bucketing; if so it falls through to a plain map join. Otherwise it compares two network costs: DPHJ moves both sides once, small plus big bytes, while a bucket-scaled map join sends the small sides to every node, nodes times small bytes. The cheaper one wins.
That comparison flips with cluster size. Suppose the small side is estimated at 3 GB and the big side at 400 GB. On 40 nodes the bucket map join costs 120 GB against 403 GB for DPHJ, so the bucket map join wins. On 200 nodes it costs 600 GB, and DPHJ wins. Same query, same data, different plan after a cluster resize.
Worked example: one join, four outcomes
Take orders, estimated at 900 GB, joined to customers, estimated at 2.4 GB with 30 million rows, under a 1 GB budget. Not bucketed, DPHJ disabled: the map join fails on size, SMB is impossible, and the result is a merge join that sorts 900 GB. Not bucketed, DPHJ enabled, about 400 reducers: 2.4 GB divided by 400 is roughly 6 MB per reducer, so DPHJ shuffles the same bytes without the sort.
Both tables bucketed into 64 buckets on customer_id, not on LLAP: 2.4 GB divided by 64 is about 37.5 MB, within budget, so the result is a bucket map join and orders is never shuffled. Filter customers to one country, dropping the estimate to 300 MB and 4 million rows, and it becomes a broadcast map join. Last, a narrow 25-million-row key table of 400 MB joined to a 6 GB fact: bytes fit, but entries exceed 21 million and the big side is under the shuffle limit, so the map join is refused; DPHJ if enabled, else a merge join.
Proving which strategy you got
Do not infer the strategy from settings; read it from the plan. Edges and operator flags identify each one.
-- Schematic EXPLAIN fragments; layout varies by release. What to look for:
Vertex dependency in root stage
Map 1 <- Map 2 (BROADCAST_EDGE) -- map join
Map 1 <- Map 2 (CUSTOM_EDGE) -- bucket map join
Reducer 3 <- Map 1 (CUSTOM_SIMPLE_EDGE), Map 2 (CUSTOM_SIMPLE_EDGE) -- DPHJ
Reducer 3 <- Map 1 (SIMPLE_EDGE), Map 2 (SIMPLE_EDGE) -- merge join
Map Join Operator -- inside a Map vertex: map join or bucket map join
DynamicPartitionHashJoin: true -- inside a Reducer: DPHJ
HybridGraceHashJoin: true -- hash table may spill partitions to disk
Merge Join Operator -- sorted shuffle join, or SMB in a Map vertex
-- BucketMapJoin: true appears at EXPLAIN EXTENDED / USER level, not the default level.A Map Join Operator in a Reducer vertex with DynamicPartitionHashJoin set is a DPHJ; a Merge Join Operator in a Map vertex is SMB. EXPLAIN ANALYZE and the Tez UI show actual rows beside estimates, exposing the estimate behind a wrong choice; Hive on Tez execution explains reading vertex metrics.
Steering the choice
Hints are weak: hive.ignore.mapjoin.hint defaults to true. The reliable levers are statistics, the budget and per-session strategy switches.
-- What is this session actually using? Defaults differ by release and vendor.
SET hive.auto.convert.join; -- master default: true
SET hive.auto.convert.join.noconditionaltask.size; -- master default: 10000000 bytes
SET hive.auto.convert.join.hashtable.max.entries; -- master default: 21000000
SET hive.auto.convert.join.shuffle.max.size; -- master default: 10000000000 bytes
SET hive.optimize.dynamic.partition.hashjoin; -- master default: false
SET hive.convert.join.bucket.mapjoin.tez; -- master default: true
SET hive.auto.convert.sortmerge.join; -- master default: true
SET hive.mapjoin.hybridgrace.hashtable; -- master default: false
SET hive.llap.mapjoin.memory.oversubscribe.factor; -- master default: 0.2
-- Every branch reads statistics. Refresh them on the tables in the join.
ANALYZE TABLE customers COMPUTE STATISTICS;
ANALYZE TABLE customers COMPUTE STATISTICS FOR COLUMNS;
-- Steer one query, then prove the plan changed.
SET hive.optimize.dynamic.partition.hashjoin=true; -- hash join instead of a sorted shuffle
SET hive.auto.convert.join.noconditionaltask.size=536870912; -- 512 MB build budget: check heap first
EXPLAIN
SELECT o.order_id, c.segment
FROM orders o JOIN customers c ON o.customer_id = c.customer_id
WHERE o.order_date >= '2026-09-01';Raise the budget only with container or executor memory to back it. Enable DPHJ where large joins sort huge inputs. hive.mapjoin.hybridgrace.hashtable lets a map join's hash table spill partitions to disk rather than fail, a safety net for underestimated build sides. Surviving shuffle joins can still be cheaper: semijoin reduction, on by default through hive.tez.dynamic.semijoin.reduction, sends a Bloom filter from the small side to prune the big side before the shuffle. Skewed keys need a different fix, described in Hive skew join optimization.
Failure modes
- Out of memory in a map join: statistics underestimated the build side; refresh column statistics, lower the budget or enable hybrid grace.
- Unexpected merge join: one input had no statistics, so map join conversion returned nothing; check that ANALYZE covered every table and partition.
- Plan changes after a cluster resize: the LLAP network-cost comparison depends on node count.
- Plan changes after a release upgrade: defaults for the budget, DPHJ or full outer handling differ between releases and vendors.
- Bucketed tables still shuffle: bucket counts or keys are incompatible, or the bucket-scaled attempt failed and the plain one was tried.
- FULL OUTER JOIN never broadcasts: in the master source read for this page, a full outer join that reaches the map join path falls back to the reduce side, where only DPHJ can handle it.
Trade-offs
A broadcast map join is the fastest plan and the most fragile, betting the task heap on an estimate. DPHJ removes the sort but still moves both inputs. Bucket map and SMB joins are cheap at query time and costly to maintain. The merge join is slow and always works. Your job is less to override the procedure than to give it accurate inputs: fresh statistics, a budget that matches real memory and a data layout that supports the strategies you want.
What to do next
- Print the join settings listed above in your environment and record the actual defaults.
- Run ANALYZE with column statistics on every table that appears in large joins, and schedule it after loads.
- For your five most expensive joins, read EXPLAIN and classify each by edge type and operator flag.
- Where a merge join sorts a large input against a medium one, test DPHJ for that session and compare runtime.
- Size noconditionaltask.size from container or executor memory, not from file sizes, and change it with headroom.
- On LLAP, re-check bucketed join plans after any change in node count.
- Add EXPLAIN output for critical queries to code review so a strategy change is visible before it ships.