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.

Advertisement

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.

FormExampleHistory
FROM (derived table)SELECT * FROM (SELECT ...) tLong supported; the subquery must have an alias. The optional AS keyword arrived in 0.13.
WHERE IN / NOT INWHERE k IN (SELECT k FROM b)0.13 (HIVE-784), single column only
WHERE EXISTS / NOT EXISTSWHERE EXISTS (SELECT 1 FROM b WHERE b.k = a.k)0.13; originally had to be correlated
Scalar in WHERE / HAVINGWHERE 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 SELECTSELECT id, (SELECT max(ts) FROM e WHERE e.id = u.id) FROM uHIVE-16091, fixed in 2.3.0; top-level expressions only
CTEWITH x AS (SELECT ...) SELECT ... FROM x0.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.

How Hive turns subqueries into joinsWHERE k IN (SELECT ...)uncorrelated or correlatedWHERE EXISTS (SELECT ...)correlated on a.k = b.kWHERE k NOT IN (SELECT ...)NULL-sensitiveWHERE NOT EXISTS (...)correlatedScalar: x > (SELECT avg ...)one row, one columnSemi joineach left row at most onceJoin + null-count checkany NULL on the right: no rowsAnti joinouter join, keep unmatchedAggregate + left joinruntime error if over one rowSince 2.2.0 (HIVE-15192) the rewrite is done by Calcite on the logical plan, not on the AST.
Each subquery form maps to a join shape; the shape decides cost and NULL behaviour.
  • 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.
Advertisement

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

SymptomCauseFix
Query returns zero rows, no errorNOT IN over a column containing NULLUse NOT EXISTS, or filter IS NOT NULL inside
Duplicated outer rowsIN rewritten by hand as an inner joinKeep IN or EXISTS (semi join), or DISTINCT the inner side
Fails hours in with a scalar subquery errorInner query returned more than one row for some keyAggregate inside the subquery
SemanticException on compileShape not supported on this version (nested WITH, complex correlation)Hoist WITH to the top; rewrite as explicit join
Same big table scanned many timesCTE inlined per reference, or several scalar subqueriesMerge into one grouped CTE, materialize, or a temp table
Incomplete output with several materialized CTEsHIVE-24606 on pre-4.0 buildsUpgrade, or use explicit temporary tables

What to do next

  1. Grep your HQL for NOT IN and replace each with NOT EXISTS unless the column is provably non-null.
  2. Run EXPLAIN on your three most expensive queries with subqueries and note the join type each became.
  3. Count TableScan operators per table; merge repeated scans into one CTE.
  4. Read hive.optimize.cte.materialize.threshold and the full.aggregate.only flag on every cluster you run on.
  5. Test materializing your heaviest multi-reference CTE, and compare wall clock and bytes read.
  6. Wrap correlated scalar subqueries in an aggregate unless the key is unique by design.
  7. If you run 2.3.x or 3.1.x and use multi-stage materialized CTEs, move them to temporary tables.
Key takeaway: Hive never executes a subquery as such: the planner rewrites IN and EXISTS into semi joins, NOT EXISTS into anti joins, NOT IN into an anti join with NULL bookkeeping, and scalar subqueries into an aggregate plus a left join checked for one row at run time. Write NOT EXISTS instead of NOT IN, aggregate inside scalar subqueries, keep WITH at the top of the statement, and remember a CTE is inlined per reference unless the materialization threshold says otherwise, a setting whose default you should read rather than assume.