Hive on Tez has a few hundred settings, and most tuning advice is a list of them with suggested values. That approach fails quietly: a setting that cured one query makes the next one slower, and nobody can say why. This article teaches a loop instead. You read what the query actually did from its plan and counters, find the stage that dominates, change the one setting family that controls that stage, and measure again.

We assume you know what a Tez DAG is. The mechanics of vertices, edges, the application master and container reuse are covered in Hive Tez execution; here we only use them to explain why a knob works. Defaults quoted below are common ones, but Hive versions and vendor distributions override several, so check yours with SET name; in Beeline before reasoning from a number.

Advertisement

Where a Hive-on-Tez query spends its time

A Hive-on-Tez query spends wall-clock time in roughly six places. Before the first task runs, a Tez application master must exist and must obtain containers from YARN. The map vertex then reads input splits that Tez has grouped into tasks. Map tasks filter, project, maybe probe a broadcast hash table, and write output that is sorted and partitioned for the next vertex. That output crosses a shuffle edge to reducers, which merge, aggregate or join, and write either to the next vertex or to the table's final location.

Each of those stages has a separate failure signature and a separate knob, which is what the diagram shows. The most important lesson is in its bottom-left corner: the cheapest work is work that never happens, so pruning and statistics come before parallelism and memory.

Session / AMprewarm, reuseSplit groupingmin/max size, wavesMap vertexscan, filter, map joinSort and spillio.sort.mb, heapShuffle edgeslow start fractionsReduce vertexbytes per reducerPruningDPP, semijoin bloomCountersspills, GC, recordsoutputfewer splitsmeasurefirst lever:less data readEach box is a cost centre with its own knob; counters tell you which one to turn
Where a Hive-on-Tez query spends time, and the setting family that controls each stage. Pruning feeds back into the scan by removing splits before they are ever grouped.

Read before you tune: plans and counters

Tuning starts with evidence. Two things give it to you cheaply: the plan and the per-vertex summary. The plan tells you which join algorithm was chosen, whether vectorization is on for each vertex, and whether dynamic partition pruning was inserted. The summary, printed after the query when hive.tez.exec.print.summary is true, shows per-vertex task counts, durations and record counts.

SET hive.tez.exec.print.summary=true;
EXPLAIN VECTORIZATION DETAIL
SELECT s.store_id, SUM(f.amount)
FROM   sales_fact f JOIN stores s ON f.store_id = s.store_id
WHERE  f.sale_date >= '2026-09-01' AND s.region = 'EMEA'
GROUP  BY s.store_id;
-- then run the query itself and read the DAG summary and counters

Then read the counters. Tez exposes them per vertex in the summary, in the Tez UI and over the timeline server. The handful below explain most slow queries:

Counter (vertex level)What a bad value looks likeUsual cause
HDFS_BYTES_READFar more than the partitions you meant to readNo partition pruning, missing ORC stats, wide SELECT
SPILLED_RECORDSGreater than output records on the map sideSort buffer too small for the map output
ADDITIONAL_SPILLS_BYTES_WRITTENNon-zero and largeMultiple spill passes, extra merge I/O
GC_TIME_MILLISMore than about 10% of task CPU timeHeap too small, or a broadcast table larger than estimated
SHUFFLE_BYTESClose to the bytes readNo map-side aggregation, filter applied too late
Task duration spreadOne task many times longer than the medianKey skew, or one oversized split

Write down the wall-clock split between vertices before you touch anything. If ninety percent of the time is in one reducer, no amount of mapper tuning will matter.

Advertisement

Startup latency: sessions, reuse and prewarm

For short interactive queries, the fixed costs dominate: launching an application master, waiting for YARN to grant containers and starting JVMs. Three things cut them. Container reuse, controlled by tez.am.container.reuse.enabled and normally on, lets a finished task's JVM run the next task. HiveServer2 session pools, configured with hive.server2.tez.initialize.default.sessions, hive.server2.tez.default.queues and hive.server2.tez.sessions.per.default.queue, keep application masters alive between queries so users do not pay for AM startup. Prewarming with hive.prewarm.enabled and hive.prewarm.numcontainers asks for containers before the first DAG arrives.

Each of these trades cluster capacity for latency. Idle sessions hold an AM container and prewarmed containers hold memory that batch jobs cannot use. Size pools to concurrent interactive users, not to total users, and let tez.session.am.dag.submit.timeout.secs reclaim abandoned sessions. If your clusters run LLAP, that is a different and larger answer to the same problem, discussed in Hive LLAP performance.

Map parallelism: the split-grouping arithmetic

The number of map tasks is not the number of files. Hive generates input splits and Tez groups them into tasks. The grouper looks at how many tasks the queue can run at once, multiplies by tez.grouping.split-waves (commonly 1.7), and divides the total input size by that to get a target group size. The target is then clamped between tez.grouping.min-size (commonly 50 MB) and tez.grouping.max-size (commonly 1 GB).

Work one example. A query reads 600 GB, and the queue can hold 400 concurrent containers. The desired task count is 400 x 1.7 = 680, so the target group is about 600 GB / 680, roughly 880 MB, inside the clamp. You get about 680 map tasks running in a little under two waves. The 0.7 extra wave is deliberate: tasks finish at different times, so a partial second wave fills gaps rather than leaving the cluster idle behind one straggler.

Now shrink the input to 8 GB. The target would be about 12 MB, below the minimum, so the clamp raises it to 50 MB and you get about 160 tasks. Each task pays JVM and setup overhead to read only 50 MB, which is why lowering tez.grouping.min-size rarely helps small queries. The opposite problem, a few very large tasks, shows up when the queue is nearly full at planning time and the grouper sees little headroom. If the input is thousands of tiny files, grouping hides some of the cost but not the metadata and open-file overhead; fix the files, as explained in the small-files deep dive.

Container, heap and sort-buffer memory

Memory settings interact, so change them together. hive.tez.container.size is the YARN container in megabytes for every Hive task. hive.tez.java.opts sets the heap; if you leave it unset, Tez derives a heap from tez.container.max.java.heap.fraction, commonly 0.8. The remainder covers off-heap memory, thread stacks and native buffers; squeezing it to zero gets containers killed by YARN for exceeding physical memory, which looks like random task failures.

Inside the heap, two consumers matter. The map-side sort buffer, tez.runtime.io.sort.mb, holds output before it spills; many managed clusters set it near 40 percent of the heap. The broadcast hash table for a map join must fit too. Hive converts a join to a map join when the estimated size of the small inputs is below hive.auto.convert.join.noconditionaltask.size, and that estimate comes from table statistics. Stale statistics are the classic cause of a map join that runs out of memory: the table is ten times larger than the metastore says. Run ANALYZE TABLE ... COMPUTE STATISTICS FOR COLUMNS before raising memory.

-- a consistent 8 GB profile; values are illustrative, size to your node memory
SET hive.tez.container.size=8192;          -- MB per task container
SET hive.tez.java.opts=-Xmx6554m;          -- about 80% of the container
SET tez.runtime.io.sort.mb=2048;           -- about 30-40% of the heap
SET tez.runtime.unordered.output.buffer.size-mb=800;  -- for unsorted edges
SET hive.auto.convert.join.noconditionaltask.size=2147483648;  -- about a third of the container, in bytes

Bigger containers are not free. Doubling container size halves how many tasks fit on a node, so a query that was CPU-bound gets slower. Raise memory only when the counters show spills or garbage collection, and prefer fixing the estimate that caused them.

Reducer count and slow start

Reducer count starts as an estimate: the expected bytes entering the reduce vertex divided by hive.exec.reducers.bytes.per.reducer (commonly 256 MB), capped by hive.exec.reducers.max. With hive.tez.auto.reducer.parallelism enabled, Hive plans for the estimate multiplied by hive.tez.max.partition.factor and lets Tez lower the count at runtime, as low as the estimate multiplied by hive.tez.min.partition.factor, once it sees the real map output sizes. It can only reduce parallelism, never raise it, so a badly low estimate stays low.

Slow start is the other reducer lever. Tez starts scheduling reducers when a fraction of source tasks has finished, between tez.shuffle-vertex-manager.min-src-fraction and tez.shuffle-vertex-manager.max-src-fraction. Starting early overlaps the shuffle fetch with the tail of the map phase. Starting too early parks reducer containers that wait for data while mappers queue for slots. On a busy shared queue, raise both fractions; on a dedicated one, the defaults usually win.

If one reducer runs far longer than the rest, the count is not your problem. That is skew, and it needs the techniques in Hive skew join optimization.

Read less: pruning, semijoins and vectorization

The biggest wins do not come from parallelism. They come from reading less. Four features do that, and each is visible in EXPLAIN.

  • Static partition pruning: filters on partition columns with literal values remove partitions at compile time. Wrapping the column in a function can stop it working.
  • Dynamic partition pruning (hive.tez.dynamic.partition.pruning): when the partition filter comes from a joined dimension, Tez sends the qualifying values to the fact scan at runtime, which skips partitions before splits are generated.
  • Semijoin reduction (hive.tez.dynamic.semijoin.reduction): for non-partition join keys, Hive builds a min/max range and a bloom filter from the small side and applies them inside the ORC reader on the big side, so non-matching rows never leave the scan.
  • Vectorization (hive.vectorized.execution.enabled): processes batches of rows per operator call. Check that EXPLAIN VECTORIZATION reports every vertex as vectorized; one unsupported UDF can silently disable it.

Join algorithm choice belongs here too, because a broadcast join removes a shuffle entirely. The decision procedure is covered in Hive join strategies.

Worked example: a 38-minute report

Here is one tuning pass. The numbers are illustrative but the sequence is typical. A daily report joins a 2 TB sales fact, partitioned by day, to a store dimension filtered to one region, and takes 38 minutes.

  1. Read. The summary shows the map vertex reading 2 TB, all partitions, even though the query asks for September. The filter is written as to_date(sale_ts) >= '2026-09-01' on a timestamp column, not on the partition column.
  2. Fix the predicate. Rewriting it against sale_date, the partition column, drops the scan to 60 GB. Runtime: 9 minutes.
  3. Read again. Now the plan shows a shuffle join. The store table is small, but its statistics are a year old and claim 4 GB. After ANALYZE TABLE stores COMPUTE STATISTICS FOR COLUMNS, Hive chooses a map join and dynamic partition pruning applies. Runtime: 4 minutes.
  4. Check memory. SPILLED_RECORDS on the map vertex is three times its output records. Raising tez.runtime.io.sort.mb within the existing heap brings spills down to one pass. Runtime: 3 minutes 10 seconds.
  5. Stop. The remaining time is split evenly across tasks with low GC. Further changes would trade cluster capacity for seconds.

Notice what we did not do: we never touched reducer counts or container sizes, which is where most people start.

Failure modes and trade-offs

Failure modeSymptomFix
Container killed for physical memoryTask attempts fail with exit code 143 or a YARN memory messageLeave about 20% of the container outside the heap; do not set -Xmx equal to the container size
Map join out of memoryOutOfMemoryError in the hash table loaderRefresh statistics; lower the conversion threshold for that query
Too many tiny tasksThousands of tasks each running a few secondsRaise the grouping minimum or compact files
Reducers starve mappersReducer containers waiting on shuffle while maps queueRaise the slow-start fractions
Settings driftA session-level SET from a script leaks into other queries in a pooled sessionPut per-query settings in the query file, not the session init

The trade-off running through all of these is cluster-wide. Every setting you raise for one query takes capacity from others on the same queue, so record why each non-default value exists, and remove it when the reason goes away.

What to do next

  1. Turn on hive.tez.exec.print.summary for your five slowest scheduled queries and record time per vertex.
  2. For each, check HDFS_BYTES_READ against the partitions you expect to read; fix predicates that defeat pruning first.
  3. Recompute column statistics on dimension tables used in joins, then re-read the plan for map joins and dynamic pruning.
  4. Confirm every vertex is vectorized with EXPLAIN VECTORIZATION.
  5. Only then compare spills and GC time against output records, and adjust sort buffer and container size together.
  6. Check grouping maths for your typical input sizes and queue capacity, and set slow start to match how busy the queue is.
  7. Keep a short file of non-default settings with the query and reason for each.
Key takeaway: Tune Hive on Tez from evidence, not lists. Read the plan and per-vertex counters, find the stage that dominates, and fix the cheapest cause first: predicates that defeat pruning, stale statistics that block map joins, and unvectorized vertices. Then size grouping, sort buffers, containers and reducers together, knowing that every capacity you give one query is taken from the rest of the queue.