A subquery is a query used as a value inside another query: as a table in FROM, as a set in WHERE ... IN, as a test in EXISTS, or as a single value. A common table expression (CTE) is a named subquery declared up front with WITH. Both make SQL easier to read. Neither exists at run time: Hive has no operator that executes a subquery. The planner rewrites every one of them into joins and aggregations before Tez ever sees the query.
That one fact explains almost everything about using them well. The rewrite decides whether a query is fast or quadratic, whether a NULL silently empties the result, and whether a CTE referenced three times is computed once or three times. This article covers what Hive accepts, what it turns each form into, the semantic traps, and how to control CTE evaluation. Join execution itself is covered in Hive Joins, in depth.
Where Hive accepts subqueries, and since when
Hive's subquery support grew in stages, and old restrictions still appear in blog posts and error messages, so it helps to know the history.
| Form | Example | History |
|---|---|---|
| FROM (derived table) | SELECT * FROM (SELECT ...) t | Long supported; the subquery must have an alias. The optional AS keyword arrived in 0.13. |
| WHERE IN / NOT IN | WHERE k IN (SELECT k FROM b) | 0.13 (HIVE-784), single column only |
| WHERE EXISTS / NOT EXISTS | WHERE EXISTS (SELECT 1 FROM b WHERE b.k = a.k) | 0.13; originally had to be correlated |
| Scalar in WHERE / HAVING | WHERE amt > (SELECT avg(amt) FROM t) | Came with the Calcite redesign (HIVE-15192, 2.2.0) and follow-up work under HIVE-15456; test on your version |
| Scalar in SELECT | SELECT id, (SELECT max(ts) FROM e WHERE e.id = u.id) FROM u | HIVE-16091, fixed in 2.3.0; top-level expressions only |
| CTE | WITH x AS (SELECT ...) SELECT ... FROM x | 0.13 (HIVE-1180); SELECT, INSERT, CTAS and CREATE VIEW |
The 0.13 manual listed restrictions that shaped a lot of Hive code still in production: subqueries only on the right-hand side of an expression, IN subqueries selecting one column, EXISTS needing a correlated predicate, and correlation allowed only in the subquery's WHERE. HIVE-15192 replaced the old AST-level rewrite with Calcite de-correlation and lifted some of these; the open umbrella HIVE-15456 tracks the rest. Treat the list as history and test the exact form on your version, because the error for an unsupported shape is a compile-time SemanticException, which is cheap to discover.
Two CTE rules have not moved in the manual: the WITH clause is not supported inside subquery blocks, and recursive queries are not supported. Put every WITH at the top of the statement, and model hierarchies another way (iterative jobs or a closure table).
The rewrite: every subquery becomes a join
With the cost-based optimizer on (see Hive CBO architecture), a subquery is parsed into a Calcite expression and a removal rule replaces it with relational operators. The shapes are worth memorising because they are what EXPLAIN will show you.
- IN and EXISTS become a semi join: each outer row is kept at most once, however many inner rows match. That is different from writing an inner join by hand, which duplicates outer rows.
- NOT EXISTS becomes an anti join, usually a left outer join followed by a filter that keeps rows where the inner side is NULL.
- NOT IN is an anti join plus extra work to honour NULL semantics, explained next.
- Scalar subqueries become an aggregate joined back with a left outer join. Hive must also prove at run time that each probe produced at most one row, and fails the query if not.
- Correlated subqueries are de-correlated: the correlation columns become join keys and the inner query is grouped by them. An equality correlation de-correlates cleanly; a non-equality one (
b.ts < a.ts) can only become a non-equi join, which is expensive or rejected.
NOT IN and NULL: the trap that returns nothing
SQL uses three-valued logic. x NOT IN (1, 2, NULL) means x <> 1 AND x <> 2 AND x <> NULL, and x <> NULL is UNKNOWN, never TRUE. So if the subquery returns even one NULL, NOT IN is never true and the query returns zero rows. No error, no warning: an empty result that looks plausible.
-- customers who never ordered: WRONG if orders.customer_id has any NULL
SELECT c.id FROM customers c
WHERE c.id NOT IN (SELECT o.customer_id FROM orders o);
-- correct and usually cheaper: NOT EXISTS ignores NULL keys
SELECT c.id FROM customers c
WHERE NOT EXISTS (SELECT 1 FROM orders o WHERE o.customer_id = c.id);To get NOT IN semantics right, the planner has to know whether the inner side contains a NULL, which means computing a count of rows and a count of non-null keys and joining that back in. That is an extra aggregation and join over the inner table. NOT EXISTS needs none of it. Rule of thumb: write NOT EXISTS by default, and use NOT IN only on a column declared or filtered NOT NULL. The join-side view of the same problem is in the joins article linked above.
Scalar subqueries
A scalar subquery must return one column and at most one row. Uncorrelated ones are easy: WHERE amount > (SELECT avg(amount) FROM sales) computes one value and cross-joins it, which is cheap because one side has a single row.
Correlated scalar subqueries in SELECT are where cost hides. (SELECT max(ts) FROM events e WHERE e.user_id = u.id) reads naturally as a per-row lookup, but Hive runs it as a full aggregation of events grouped by user_id, then a left join. That is often fine, and is exactly what you would write by hand. It becomes a problem when a query has five such columns over the same table: five aggregations and five joins, where one grouped CTE would do. If a scalar subquery can return two rows for some key (no aggregate, duplicate keys), the query fails at run time, possibly hours in, so wrap the inner select in an aggregate whenever uniqueness is not guaranteed by the data model.
CTEs: scope, inlining and materialization
A CTE is scoped to one statement. By default Hive treats it like a view: each reference is replaced by the CTE's query text. A CTE referenced three times is therefore computed three times unless the optimizer happens to share the work. For a cheap filter that is fine; for a CTE that scans a large table and aggregates, it triples the cost.
Hive can instead materialize a CTE into a temporary table before the main query runs. Two settings control it: hive.optimize.cte.materialize.threshold (materialize a CTE referenced at least this many times; -1 disables) and hive.optimize.cte.materialize.full.aggregate.only (restrict materialization to CTEs whose output is a full aggregate). Defaults have differed between Hive releases and vendor builds, so do not assume; read them on your cluster.
SET hive.optimize.cte.materialize.threshold; -- print the current value
SET hive.optimize.cte.materialize.full.aggregate.only;
SET hive.optimize.cte.materialize.threshold=2; -- for this session only
SET hive.optimize.cte.materialize.full.aggregate.only=false;Materialization has costs too: a write to HDFS and a read back, and a barrier between stages that removes pipelining. It wins when the CTE is expensive and its output is small. It also had a real correctness bug: HIVE-24606 reports that with several materialized CTEs plus a non-materialized one, the dependency between materialized CTEs could be lost and a later CTE could run before an earlier one finished, producing incomplete output. It affects 2.3.7 and 3.1.2 and is fixed in 4.0.0. On older versions, prefer an explicit temporary table for multi-stage pipelines.
CREATE TEMPORARY TABLE active_30d STORED AS ORC AS
SELECT customer_id, count(*) AS orders, sum(amount) AS spend
FROM orders
WHERE order_date >= date_sub(current_date, 30)
GROUP BY customer_id;
-- reuse active_30d in as many statements as needed; it is dropped at session end
Worked example: churn-risk customers
Task: list customers who ordered in the previous 90 days but not in the last 30, with their top category by spend and their spend relative to the segment average. A first draft piles subqueries into the select list.
SELECT c.id,
(SELECT o.category FROM orders o WHERE o.customer_id = c.id
ORDER BY o.amount DESC LIMIT 1) AS top_cat, -- rejected or slow
(SELECT sum(o.amount) FROM orders o WHERE o.customer_id = c.id) AS spend
FROM customers c
WHERE c.id IN (SELECT customer_id FROM orders
WHERE order_date >= date_sub(current_date, 90))
AND c.id NOT IN (SELECT customer_id FROM orders
WHERE order_date >= date_sub(current_date, 30));It has three problems: a correlated subquery with ORDER BY ... LIMIT is not a top-level aggregate and is either rejected or planned badly; orders is scanned four times; and NOT IN returns nothing if any recent order has a NULL customer id. The rewrite reads orders once into a CTE, ranks with a window function (see Hive window functions), and uses NOT EXISTS.
WITH recent AS (
SELECT customer_id, category, amount, order_date
FROM orders
WHERE order_date >= date_sub(current_date, 90) AND customer_id IS NOT NULL
),
cat_spend AS (
SELECT customer_id, category, sum(amount) AS spend
FROM recent GROUP BY customer_id, category
),
ranked AS (
SELECT customer_id, category, spend,
row_number() OVER (PARTITION BY customer_id ORDER BY spend DESC) AS rn
FROM cat_spend
),
totals AS (
SELECT customer_id, sum(spend) AS spend FROM cat_spend GROUP BY customer_id
),
seg_avg AS ( -- all active customers, before the churn filter
SELECT c.segment, avg(t.spend) AS avg_spend
FROM customers c JOIN totals t ON t.customer_id = c.id
GROUP BY c.segment
)
SELECT c.id, c.segment, r.category AS top_cat, t.spend,
t.spend / s.avg_spend AS vs_segment
FROM customers c
JOIN totals t ON t.customer_id = c.id
JOIN ranked r ON r.customer_id = c.id AND r.rn = 1
JOIN seg_avg s ON s.segment = c.segment
WHERE NOT EXISTS (SELECT 1 FROM recent r
WHERE r.customer_id = c.id
AND r.order_date >= date_sub(current_date, 30));Run EXPLAIN on both. Aggregation and ranking are split into two CTEs so the window never wraps an aggregate, and the segment average is computed in its own CTE: a window average in the final select would run after the WHERE and average only the at-risk customers. recent is referenced twice (through cat_spend and the NOT EXISTS), and totals twice, so with a threshold of 2 they are candidates for materialization; whether that wins depends on how selective the 90-day filter is. The JOIN totals already restricts to customers active in 90 days, so the original IN subquery disappears entirely.
Reading the plan
Three things to look for in EXPLAIN output for any query with subqueries. First, the join type that replaced each subquery: a left semi join for IN or EXISTS, an outer join with an IS NULL filter for an anti join. Second, how many times each large table is scanned: count the TableScan operators per alias. Third, whether a semi join became a map join: if the inner side is small after filtering, a broadcast join avoids a shuffle entirely. If the plan you expected did not appear, check that the CBO actually ran, because statements that fall back to the legacy path get the old rewrite. The planner's statistics, which drive all of this, are covered in Hive Cost-Based Optimizer.
Failure modes
| Symptom | Cause | Fix |
|---|---|---|
| Query returns zero rows, no error | NOT IN over a column containing NULL | Use NOT EXISTS, or filter IS NOT NULL inside |
| Duplicated outer rows | IN rewritten by hand as an inner join | Keep IN or EXISTS (semi join), or DISTINCT the inner side |
| Fails hours in with a scalar subquery error | Inner query returned more than one row for some key | Aggregate inside the subquery |
| SemanticException on compile | Shape not supported on this version (nested WITH, complex correlation) | Hoist WITH to the top; rewrite as explicit join |
| Same big table scanned many times | CTE inlined per reference, or several scalar subqueries | Merge into one grouped CTE, materialize, or a temp table |
| Incomplete output with several materialized CTEs | HIVE-24606 on pre-4.0 builds | Upgrade, or use explicit temporary tables |
What to do next
- Grep your HQL for NOT IN and replace each with NOT EXISTS unless the column is provably non-null.
- Run EXPLAIN on your three most expensive queries with subqueries and note the join type each became.
- Count TableScan operators per table; merge repeated scans into one CTE.
- Read hive.optimize.cte.materialize.threshold and the full.aggregate.only flag on every cluster you run on.
- Test materializing your heaviest multi-reference CTE, and compare wall clock and bytes read.
- Wrap correlated scalar subqueries in an aggregate unless the key is unique by design.
- If you run 2.3.x or 3.1.x and use multi-stage materialized CTEs, move them to temporary tables.