A join answers two questions, and most Hive join problems come from mixing them up. The logical one: which rows should come back, given the join type, ON, WHERE and NULLs? The physical one: does the engine shuffle both tables, broadcast a small one, or exploit the data layout? A query can be fast and wrong, or right and slow, and the fixes differ completely.
This article covers semantics first, including the filter-placement and NULL-key rules that silently change results, then the execution strategies, the settings and statistics that choose between them, and how to confirm the choice in EXPLAIN. Bucketed joins, skew and cost-based join ordering have their own articles, linked where they come up.
Join types and what each returns
Hive supports the standard SQL join types plus one of its own. Think of each in terms of which side is preserved, meaning its rows survive even without a match, and how many times a row can appear.
| Join | Rows returned | Typical use |
|---|---|---|
| INNER JOIN | Only pairs that satisfy ON; a row with k matches appears k times | Facts with required dimensions |
| LEFT OUTER JOIN | Every left row; right columns NULL when unmatched | Enrich facts with optional attributes |
| FULL OUTER JOIN | Every row from both sides | Reconciliation, diffing two snapshots |
| LEFT SEMI JOIN | Left rows that have at least one match, each once; right columns not selectable | Existence filters |
| CROSS JOIN | Every combination; rows multiply | Small calendars or parameter grids only |
The multiplicity column matters more than people expect. An inner join does not filter a fact table, it multiplies it. If a dimension that should have one row per key has two, every matching fact row doubles, and sums computed after the join are silently wrong. Hive does not check key uniqueness for you.
ON versus WHERE in outer joins
In an inner join it makes no difference whether a filter is written in ON or in WHERE. In an outer join it changes the answer. ON decides which pairs match. WHERE runs afterwards, on the rows the join produced, including the NULL-extended ones. Three versions of the same query show the effect:
-- customers orders
-- id | region order_id | cust_id | status
-- 1 | EU 10 | 1 | PAID
-- 2 | US 11 | 1 | CANCELLED
-- 3 | EU 12 | 2 | CANCELLED
-- A: filter in ON. Every customer survives; only PAID orders attach.
SELECT c.id, o.order_id
FROM customers c
LEFT JOIN orders o ON o.cust_id = c.id AND o.status = 'PAID';
-- 1 | 10
-- 2 | NULL
-- 3 | NULL
-- B: filter in WHERE. Runs after the join, removes NULL-extended rows: now an inner join.
SELECT c.id, o.order_id
FROM customers c
LEFT JOIN orders o ON o.cust_id = c.id
WHERE o.status = 'PAID';
-- 1 | 10
-- C: filter on the preserved side in ON does NOT remove customers.
SELECT c.id, o.order_id
FROM customers c
LEFT JOIN orders o ON o.cust_id = c.id AND c.region = 'EU';
-- 1 | 10
-- 1 | 11
-- 2 | NULL <- still returned, just never matched
-- 3 | NULLQuery A keeps every customer and attaches only paid orders. Query B filters the NULL-producing side in WHERE, which rejects every NULL-extended row and turns the outer join into an inner join. Query C is the surprise: a preserved-side condition in ON removes no customers, it only stops them matching. Rule: NULL-producing-side conditions go in ON, preserved-side conditions go in WHERE.
The optimizer may plan B as an inner join, since its WHERE rejects NULLs anyway, so the plan for a buggy query looks perfectly healthy.
NULL keys, semi joins and anti joins
Equality with NULL is never true, so rows with a NULL join key never match in an equi-join. In an inner join they vanish. In a left outer join they survive with NULL right columns, which is right but often unnoticed. When NULL is a meaningful value, such as a missing tenant that should match other missing tenants, use the null-safe operator <=>, which treats two NULLs as equal, aware that many NULL keys matching each other can explode the join.
-- Customers with at least one order (each customer once, however many orders).
SELECT c.id FROM customers c LEFT SEMI JOIN orders o ON o.cust_id = c.id;
-- The same intent, written portably.
SELECT c.id FROM customers c
WHERE EXISTS (SELECT 1 FROM orders o WHERE o.cust_id = c.id);
-- Anti join: customers with no orders.
SELECT c.id FROM customers c
LEFT JOIN orders o ON o.cust_id = c.id
WHERE o.cust_id IS NULL;
-- Null-safe equality: NULL matches NULL.
SELECT * FROM a JOIN b ON a.tenant <=> b.tenant;LEFT SEMI JOIN returns each left row at most once, so it cannot multiply rows, and Hive also rewrites EXISTS and IN subqueries into semi joins. For anti joins, use LEFT JOIN with an IS NULL test on the right join key, a column that is never NULL in matched rows.
Shuffle joins: the general case
When nothing better applies, Hive runs a shuffle join, historically called a common join. Both sides are scanned, repartitioned and sorted by the join key, and join tasks merge matching ranges; on Tez this is a Merge Join Operator in a reducer vertex. It works for every join type and size, but moves both inputs across the network.
Joins on the same key are combined into one shuffle stage. Old advice to put the largest table last, or name it with the STREAMTABLE hint, comes from MapReduce, where the reducer buffered all but the last table per key. With the cost-based optimizer and statistics, join order is planned for you, as described in the Hive CBO article.
MapJoin: broadcasting the small side
If one side is small enough to fit in memory, a far cheaper plan exists. Hive reads the small table, builds a hash table keyed by the join key, and sends a copy to every task that scans the large table. The large table is never shuffled. On Tez the copy travels over a broadcast edge, and the plan shows a Map Join Operator inside the large table's map vertex.
The planner converts a shuffle join into a MapJoin when hive.auto.convert.join is on and the estimated small-side size is below a threshold; with hive.auto.convert.join.noconditionaltask, several MapJoins whose small tables together fit under hive.auto.convert.join.noconditionaltask.size are chained into one task. Estimates come from statistics, so stale statistics cause missed MapJoins or, worse, MapJoins on large tables. Defaults vary by release and vendor, so read them:
-- Read the values your cluster actually uses; defaults differ by release and vendor.
SET hive.auto.convert.join;
SET hive.auto.convert.join.noconditionaltask;
SET hive.auto.convert.join.noconditionaltask.size;
SET hive.ignore.mapjoin.hint;
SET hive.tez.container.size;
-- Statistics drive the size estimate. Without them, the planner guesses.
ANALYZE TABLE customers COMPUTE STATISTICS;
ANALYZE TABLE customers COMPUTE STATISTICS FOR COLUMNS;Not every side can be the hashed one. The broadcast side must be one whose unmatched rows never need to be emitted, because each probe task sees only its own slice of the large table and cannot know whether another task found a match.
| Join type | Side that can be hashed and broadcast | Consequence |
|---|---|---|
| INNER | Either side | Planner picks the smaller one |
| LEFT OUTER | Right side only | Small left, big right: no MapJoin |
| RIGHT OUTER | Left side only | Mirror image |
| LEFT SEMI | Right side | Natural fit; existence probe |
| FULL OUTER | Neither, classically | Shuffle join; newer releases may differ, so check the plan |
If a MAPJOIN hint seems to do nothing, check hive.ignore.mapjoin.hint. When the large table is also bucketed or sorted on the key, bucket map joins and sort-merge-bucket joins avoid even the broadcast, which is covered in bucketing for joins.
Making joins cheaper before they run
The fastest join reads less data. Dynamic partition pruning on Tez (hive.tez.dynamic.partition.pruning) sends the dimension's qualifying keys to the fact scan at run time so whole partitions are skipped. Semijoin reduction (hive.tez.dynamic.semijoin.reduction) applies a Bloom filter and min/max range from the small side to the large side's scan, much like Impala runtime filters. Add predicate pushdown into ORC and Parquet: filter early, join late.
Reading the plan
Never assume which strategy ran. EXPLAIN shows it, and the operators and edge types are what to look for:
EXPLAIN
SELECT o.order_id, c.region
FROM orders o JOIN customers c ON o.cust_id = c.id
WHERE o.dt = '2026-09-30';
-- Abridged and illustrative; operator names are what to look for.
-- Edges:
-- Map 1 <- Map 2 (BROADCAST_EDGE)
-- Map 1
-- TableScan alias: o Filter Operator predicate: cust_id is not null
-- Map Join Operator
-- condition map: Inner Join 0 to 1
-- keys: 0 cust_id (type: bigint) 1 id (type: bigint)
-- input vertices: 1 Map 2
-- Map 2
-- TableScan alias: c Filter Operator predicate: id is not null
-- Reduce Output Operator key expressions: id (feeds the broadcast)- Map Join Operator inside a map vertex, with an input vertex over a BROADCAST_EDGE: a MapJoin. The broadcast vertex is the hashed side; confirm it is the one you expect to be small.
- Merge Join Operator in a reducer vertex fed by SIMPLE_EDGE inputs: a shuffle join.
- Key types in the join keys line: both sides should be the same type. A cast in the keys line is a warning sign, as the next section explains.
At run time, a MapJoin whose build vertex reads far more rows than statistics claimed signals a coming memory failure. Hive on Tez execution explains how to read those vertices.
Worked example: one fact, two dimensions
A daily report joins orders, partitioned by date with about 40 million rows per day, to customers with 5 million rows, and to a 400-row regions table. The first run takes 14 minutes. EXPLAIN shows two Merge Join Operators: both joins are shuffling. DESCRIBE FORMATTED customers shows no column statistics and a row count of zero, because the table is rebuilt nightly by a job that never ran ANALYZE.
After computing statistics, regions becomes a MapJoin at once; customers stays above the threshold, and the run drops to about six minutes. Rather than raising the threshold, which would risk a hash table several times the on-disk size in a 4 GB container, the team projects only the two needed columns and filters to active customers before the join. That brings the estimate under the limit honestly, and the run takes under two minutes. The figures are illustrative; the order is the point: statistics, plan, measurement, thresholds last.
Failure modes
- Row explosion. A duplicate key in a supposedly unique dimension multiplies facts. Check with GROUP BY key HAVING count(*) greater than 1, and fail the pipeline on any result.
- Silent type coercion. Joining a STRING key to a BIGINT key forces a common type. With some type pairs the comparison goes through DOUBLE, which cannot represent every 64-bit integer, so distinct large IDs can compare equal and produce false matches. Cast explicitly and lossless, and count values that fail the cast:
-- orders.cust_ref is STRING, customers.id is BIGINT.
-- Hive must reconcile the types before comparing; with this pair the comparison
-- can go through DOUBLE, which only represents integers exactly up to 2^53.
SELECT ... FROM orders o JOIN customers c ON o.cust_ref = c.id; -- risky
-- Make the conversion explicit and lossless, and check the cast's failures.
SELECT ... FROM orders o JOIN customers c ON CAST(o.cust_ref AS BIGINT) = c.id;- MapJoin out of memory. Statistics said small, data says large. Refresh statistics after bulk loads and treat the failure as a planning bug, not a capacity problem.
- Accidental inner join. A WHERE filter on the NULL-producing side of an outer join, as in query B above. Code review should check every outer join's filters.
- Cartesian products. A missing ON condition yields a cross product. Keep your version's strict-mode check on in shared clusters.
- Skewed keys. One key with a large share of rows sends all of them to one join task, which runs for hours while the rest finish. See skew join optimization for detection and fixes.
What to do next
- Pick your three most expensive join queries and run EXPLAIN on each. Write down, for every join, whether it is a MapJoin or a shuffle join and which side is hashed.
- Run the SET commands above on your cluster and record the real values for conversion, thresholds and hint handling.
- Check that every table involved has fresh table and column statistics, and add ANALYZE to the jobs that rebuild them.
- Review every outer join for filters in the wrong clause, using the ON versus WHERE rule.
- Add a uniqueness check for every dimension key used in joins, and a type check that both sides of each join key share a type.
- Only after all of the above, consider raising MapJoin thresholds, and measure memory use when you do.