Hive on Tez turns one SQL statement into one directed acyclic graph of vertices, and the speed of the query is decided by a handful of numbers: how many tasks each vertex runs, how much memory each task gets, how data moves along each edge, and how much input is skipped before any task starts. Most slow Hive queries are not slow because Tez is slow. They are slow because one of those numbers was left to a default that does not fit the data.

Our Hive on Tez architecture article explains the application master, sessions and container reuse. This deep dive sits one layer down: how to read the plan Hive hands to Tez, which settings control each vertex, how runtime decisions such as auto reducer parallelism and dynamic partition pruning actually work, and how to tune a real query by reading its counters instead of guessing. Settings are named as they appear in Apache Hive and Tez; distributions sometimes ship different defaults, so always check your own values with SET name; before changing them.

Advertisement

From SQL to a Tez DAG: what the edges tell you

Hive compiles SQL into an operator tree, the cost-based optimizer reorders joins using table and column statistics, and the Tez compiler cuts the tree into vertices wherever data must be redistributed. A vertex named Map N reads a table or file input; a vertex named Reducer N reads only from other vertices. The numbers are labels, not execution order. What matters is the edge between vertices, because the edge type tells you what kind of data movement you are paying for.

EXPLAIN edgeTez movementTypical causeCost to watch
SIMPLE_EDGEScatter-gather, sorted by keyGROUP BY, ORDER BY, shuffle joinSort spills, shuffle bytes, skew
BROADCAST_EDGEEvery task gets the whole outputMap join of a small tableHash table memory in every task
CUSTOM_SIMPLE_EDGEPartitioned shuffle without sortDynamically partitioned hash joinMemory for per-partition hash tables
CUSTOM_EDGEBucket-aware routingBucket map joinRequires matching bucket layout
XPROD_EDGECross product distributionCartesian joinOutput size explodes

A plan with one SIMPLE_EDGE per GROUP BY and BROADCAST_EDGE for every dimension is usually healthy. A SIMPLE_EDGE where you expected a broadcast means the planner decided a table was too big to map join, usually because statistics are missing or stale; the join strategy article covers that decision in detail.

Reading the plan before touching a setting

Start every tuning session by printing the plan and the per-vertex summary. The user-level EXPLAIN gives the vertex dependency graph at the top; EXPLAIN FORMATTED gives JSON if you want to diff plans in a script.

SET hive.tez.exec.print.summary=true;   -- per-vertex table after each query in Beeline
EXPLAIN
SELECT s.region, SUM(f.amount) AS revenue
FROM   sales f
JOIN   dim_store s ON f.store_id = s.store_id
JOIN   dim_date  d ON f.sale_date = d.date_key
WHERE  d.month = '2026-09'
GROUP  BY s.region
ORDER  BY revenue DESC;

-- Vertex dependency in root stage
-- Map 1 <- Map 4 (BROADCAST_EDGE), Map 5 (BROADCAST_EDGE)
-- Reducer 2 <- Map 1 (SIMPLE_EDGE)
-- Reducer 3 <- Reducer 2 (SIMPLE_EDGE)

Read three things in the operator detail. First, the Statistics: Num rows ... Data size ... lines: if they say basic stats: PARTIAL or show absurdly small sizes, every later decision is built on a bad estimate, so run ANALYZE TABLE ... COMPUTE STATISTICS FOR COLUMNS first. Second, look for a Dynamic Partitioning Event Operator under the dimension vertex; its absence on a partitioned fact table means no runtime pruning. Third, check that the map join operators sit inside Map 1 rather than in a reducer. The diagram below shows the shape this query should have.

Map 4scan dim_date, filterMap 5scan dim_storeMap 1scan sales, map joinReducer 2partial to final aggregateReducer 3ORDER BY, one taskSplit generator for Map 1prunes partitions, groups splitsBROADCAST_EDGEBROADCAST_EDGEDPP eventgrouped splitsSIMPLE_EDGESIMPLE_EDGEEach vertex runs N tasks; edges decide how one vertex's output partitions reach the next
A typical star-join query on Tez: two small dimension scans broadcast hash tables into the fact scan, the date dimension also sends a dynamic partition pruning event to the fact scan's split generator, and two scatter-gather shuffles feed the aggregate and the final sort.
Advertisement

Mapper parallelism comes from split grouping

A Map vertex's task count is not set directly. The split generator lists the files of every partition that survives pruning, asks the input format for splits, and then Tez groups those splits into larger units so that the task count matches the cluster rather than the file layout. The grouping target is roughly the number of containers the queue can hold multiplied by tez.grouping.split-waves (1.7 by default, so a little more than one wave of tasks), and each group is kept between tez.grouping.min-size (16 MB) and tez.grouping.max-size (1 GB).

Two consequences follow. A table made of thousands of tiny files does not produce thousands of tasks, because grouping merges them, but each task still opens every file, so open and footer-read latency dominates; fix the files with compaction or INSERT OVERWRITE rather than with grouping settings. And when a busy queue has little headroom at submission time, the generator plans fewer, larger groups, so the same query can run with very different mapper counts at different times of day. If you need predictable mapper counts for a benchmark, pin tez.grouping.min-size and tez.grouping.max-size close together for that session.

Reducer parallelism: estimate, then shrink at runtime

For a Reducer vertex Hive first estimates a count at compile time: the estimated bytes flowing into the vertex divided by hive.exec.reducers.bytes.per.reducer (256 MB in recent Apache Hive), capped by hive.exec.reducers.max. Setting mapreduce.job.reduces to a positive number overrides the estimate entirely, which is almost always a mistake left over from MapReduce scripts.

With hive.tez.auto.reducer.parallelism=true the compiler deliberately over-provisions: it plans the estimate multiplied by hive.tez.max.partition.factor (2.0) partitions, and Tez's ShuffleVertexManager watches the real output size of the upstream tasks. Once enough of them have finished, it merges adjacent partitions so that each reducer receives about the target byte count, never going below the estimate multiplied by hive.tez.min.partition.factor (0.25). The key property: runtime adjustment can only reduce the reducer count. If statistics underestimate the data by ten times, auto parallelism cannot rescue you; the vertex starts with too few partitions and each reducer spills.

Shuffle: slow-start, sort buffers and skew

Reducers do not wait for every mapper. The ShuffleVertexManager starts scheduling reducer tasks when tez.shuffle-vertex-manager.min-src-fraction of source tasks are done (0.25 by default) and has all of them running by tez.shuffle-vertex-manager.max-src-fraction (0.75). Early start overlaps fetching with mapping, which helps long map phases, but on a full queue early reducers occupy containers while they sit idle waiting for the last mappers. If the Tez UI shows reducers running for minutes with almost no CPU, raise both fractions.

On the map side, each task of a sorted edge writes into an in-memory buffer of tez.runtime.io.sort.mb; when it fills, the task sorts and spills a run to local disk and merges runs at the end. Unsorted edges use tez.runtime.unordered.output.buffer.size-mb instead. The counters tell you whether the buffers are too small: SPILLED_RECORDS much larger than output records, or a large ADDITIONAL_SPILLS_BYTES_WRITTEN, means repeated spill and merge passes. On the fetch side, SHUFFLE_BYTES per reducer task shows skew directly; one task fetching twenty times the median is a hot key, not a parallelism problem, and needs hive.optimize.skewjoin or a salted key.

Dynamic partition pruning events

Dynamic partition pruning is the most valuable runtime optimisation for star schemas. When the planner sees a join between a partition column of the fact table and a filtered dimension, it adds a branch to the dimension vertex that collects the distinct join key values and sends them, as a Tez input initializer event, to the fact vertex's split generator. The generator waits for the event, drops every partition not in the set, and only then plans splits. A query over three years of daily partitions that filters to one month reads about thirty partitions instead of more than a thousand.

The feature is controlled by hive.tez.dynamic.partition.pruning and bounded by hive.tez.dynamic.partition.pruning.max.event.size and hive.tez.dynamic.partition.pruning.max.data.size. When the planner estimates that the pruning values would exceed those limits, it removes the branch and the fact scan reads everything. It also cannot prune when the fact side of the join is an expression rather than the bare partition column, for example CAST(f.sale_date AS STRING) = d.date_str, or when the dimension filter is too unselective to be worth it. Always confirm the event operator in EXPLAIN; it is the cheapest win in this article.

Container and memory sizing

Each task runs in a YARN container of hive.tez.container.size MB, with JVM options from hive.tez.java.opts. A value of -1 for the container size falls back to the MapReduce map memory setting, which is rarely what you want on a modern cluster. The heap must be smaller than the container to leave room for metaspace, thread stacks and direct buffers, so a common rule is a heap of about 80 percent of the container; set it explicitly rather than trusting a fraction to be computed for you.

-- A session profile for a heavy reporting query on 8 GB containers
SET hive.tez.container.size=8192;
SET hive.tez.java.opts=-Xmx6554m -XX:+UseG1GC;
SET tez.runtime.io.sort.mb=1600;                         -- well under half the heap
SET hive.auto.convert.join.noconditionaltask.size=1700000000;  -- map join hash table budget, bytes
SET tez.grouping.max-size=536870912;                     -- 512 MB groups, more mappers

Three budgets share that heap: the sort buffer, the map join hash tables, and the operators themselves (aggregation hash maps, vectorized batches). hive.auto.convert.join.noconditionaltask.size is the total size of all small tables that may be combined into one map join task; set it to roughly a third of the heap so that a broadcast of several dimensions plus the sort buffer still fits. Bigger containers are not free: fewer of them fit on a node, which lowers the grouping target and the parallelism of every vertex. Size up only the queries that need it, through session settings or a resource plan, not the cluster default.

Worked example: tuning one query by its counters

Take the revenue query above on a fact table of about 2 TB of ORC, partitioned by day over three years. The numbers in this walk-through are illustrative, but the sequence is the one that works in practice. The first run took 41 minutes. The print summary showed Map 1 with 7,900 tasks reading 2.0 TB, Reducer 2 with 380 tasks, and Reducer 3 with one task finishing in seconds.

Step 1, pruning. Reading 2 TB for one month of data meant DPP was missing. EXPLAIN had no event operator because the view behind sales joined on CAST(sale_date AS STRING). Rewriting the view to expose the typed partition column restored the event; Map 1 dropped to 330 tasks reading 56 GB. Runtime: 6 minutes.

Step 2, spills. Map 1 counters showed SPILLED_RECORDS at 3.1 times output records with a 100 MB sort buffer on 4 GB containers. Moving to 6 GB containers with a 1.6 GB buffer took spills to 1.0 times, meaning a single spill per task. Runtime: 4 minutes 20 seconds.

Step 3, reducers. Reducer 2 had been planned at 380 tasks from a stale estimate, then shrunk at runtime to 95, and still most tasks finished in under three seconds, a sign of too many tiny reducers. After computing column statistics, the estimate dropped to 56, auto parallelism settled at 28, and the summary showed even task durations. Runtime: 3 minutes 40 seconds.

RunMap 1 tasks / inputSpill ratioReducer 2 tasksWall time
Baseline7,900 / 2.0 TB3.1380 planned, 95 run41 min
DPP restored330 / 56 GB3.1380 planned, 90 run6 min
Bigger sort buffer330 / 56 GB1.0380 planned, 90 run4 min 20 s
Fresh column stats330 / 56 GB1.056 planned, 28 run3 min 40 s

Notice the order: fix the bytes read first, then the per-task efficiency, then the parallelism. Tuning reducers before restoring pruning would have optimised work that should never have happened.

Failure modes and trade-offs

  • Container killed by YARN for exceeding memory. The heap plus off-heap use outgrew the container. Lower the heap fraction, or raise the container size together with the heap, never the heap alone.
  • Map join out of memory. A dimension grew past the planner's estimate. Refresh statistics, lower the no-conditional-task size, or let the planner fall back to a shuffle join.
  • A vertex stuck at 99 percent. One straggler reducer means skewed keys; check per-task shuffle bytes before adding reducers, because more partitions do not split a single hot key.
  • Queries slow only at peak hours. Split grouping planned fewer tasks because the queue had little headroom, and early-started reducers held containers. Use separate queues or LLAP for interactive work; the LLAP architecture article explains the persistent-daemon alternative.
  • Fast in testing, slow in production. Test tables had fresh statistics and production ones did not. Make ANALYZE TABLE part of the load pipeline.

The trade-offs are consistent. More parallelism shortens the critical path but adds scheduling overhead, shuffle connections and small output files. Bigger containers reduce spills and permit broadcast joins but reduce concurrency for everyone else. Aggressive slow-start helps a single long query and hurts a crowded queue. Vectorized execution, covered in our vectorization article, reduces CPU per row and is worth enabling before any of these memory changes.

What to do next

  1. Pick your three slowest recurring queries and run each with hive.tez.exec.print.summary=true and EXPLAIN; save both as a baseline.
  2. For every partitioned fact table, confirm a Dynamic Partitioning Event Operator appears in the plan; rewrite joins that wrap the partition column in an expression.
  3. Check statistics: run DESCRIBE FORMATTED on each table and partition, and add column statistics to the load job where they are missing.
  4. Compare SPILLED_RECORDS to output records per vertex; if the ratio is above about 1.5, size up the sort buffer within an explicit heap budget.
  5. Look at per-task shuffle bytes for each reducer vertex and treat a single outlier as skew, not as a parallelism setting.
  6. Write the settings that helped into a named session profile or resource plan, not the global default, and re-measure after each change, one change at a time.
Key takeaway: Hive on Tez performance is set by a few numbers per vertex: tasks, memory, edge type and bytes read. Read EXPLAIN first, confirm dynamic partition pruning on partitioned fact tables, and keep statistics fresh because reducer estimates and map join decisions depend on them and runtime auto parallelism can only shrink. Then use spill and shuffle counters to size sort buffers and spot skew, and change one setting at a time inside a session profile.