A Hive query touches at least four kinds of Java virtual machine: HiveServer2, which parses and plans the query; the Metastore, which answers questions about tables and partitions; the Tez Application Master, which schedules the DAG; and dozens to thousands of Tez task JVMs running inside YARN containers. Each has its own memory settings, its own way of failing, and its own logs. Most "Hive is slow" or "Hive OOMs" tickets come from tuning the wrong one.
This article maps those JVMs, explains how the heap relates to the YARN container around it, gives a sizing method for Tez tasks with a worked example, covers HiveServer2 and Metastore heaps and garbage collector choice, and ends with a table that maps the error messages you will actually see to the setting that fixes them. It assumes Hive on Tez (Hive 3 and 4); LLAP daemons have their own memory model, covered in LLAP performance.
The JVMs in a Hive deployment
Start by knowing which process failed, because the same words, OutOfMemoryError, mean different things in each:
| JVM | Lifetime | Memory driven by | Main settings |
|---|---|---|---|
| HiveServer2 | Long-running service | Concurrent sessions, query compilation (plans over many partitions), result fetch buffers, operation logs | Heap in hive-env.sh or the cluster manager; GC flags |
| Metastore | Long-running service | Partition listings, column statistics, object caches, concurrent Thrift calls | Heap, GC flags, backing database pool |
| Tez AM | Per session or per query | Number of tasks and splits, DAG size, task events and counters | tez.am.resource.memory.mb, tez.am.launch.cmd-opts |
| Tez task | Per container; reused across tasks | Sort and shuffle buffers, map-join hash tables, operator state, ORC/Parquet readers | hive.tez.container.size, hive.tez.java.opts, buffer sizes |
The services are tuned rarely and carefully, with restarts. Task containers can be tuned per query with SET in the session, which makes them the right place to fix one heavy query without touching everyone else.
Container versus heap
YARN grants each Tez task a container of hive.tez.container.size megabytes (if that is left at -1, Hive falls back to the MapReduce map memory setting). The NodeManager monitors the process tree's physical memory and kills it if it exceeds the grant. Inside, the JVM's heap is capped by -Xmx, but the process also uses memory the heap does not count: metaspace for classes, thread stacks, the JIT's code cache, direct NIO buffers used by compression codecs and network transfers, and native allocations by libraries such as zlib and Snappy.
So the heap must be smaller than the container. Tez handles this with tez.container.max.java.heap.fraction, which defaults to 0.8: if no -Xmx is supplied in the task options, Tez sets the heap to 80% of the container. If you set hive.tez.java.opts yourself with an explicit -Xmx, that value wins, and it is your job to keep it at roughly 75-85% of the container. Setting -Xmx equal to the container size is the most common mistake: the heap fills, the off-heap use pushes the process over, and YARN kills a JVM that never reported an OutOfMemoryError.
YARN also rounds every request up to a multiple of yarn.scheduler.minimum-allocation-mb. A 4,500 MB request on a scheduler with a 1,024 MB minimum gets 5,120 MB, so choose container sizes that are multiples of the minimum or you waste the difference on every container.
What fills a task heap
Inside the heap, three consumers dominate, and Hive and Tez size them independently, which is why changing the container alone often does not help:
- Sort buffer (
tez.runtime.io.sort.mb): the in-memory buffer for ordered output, used by shuffle edges for joins and aggregations. Too small means many spills; too large starves everything else. - Unordered output buffer (
tez.runtime.unordered.output.buffer.size-mb): used by broadcast and unsorted edges. - Map-join hash tables: when Hive converts a join to a map join, the small side is loaded into a hash table in every task.
hive.auto.convert.join.noconditionaltask.sizecaps the total estimated size of the small tables merged into one task. The in-memory table is larger than the on-disk estimate because of object overhead, and stale statistics make the estimate wrong in either direction.
Commonly quoted starting points, from vendor tuning guides rather than any specification, are a sort buffer of about 40% of the task heap, an unordered buffer of about 10%, and a map-join threshold of about one third of the heap. Some guides quote the same ratios against the container size; using the heap is the more conservative reading. Treat them as a first guess to measure against, not as rules. The Hive joins guide explains when map joins are chosen.
Worked example: sizing a 128 GB node
A worker node has 128 GB of RAM and 32 cores. After the OS, DataNode and NodeManager daemons, YARN is given 108 GB (yarn.nodemanager.resource.memory-mb=110592), with a 1,024 MB minimum allocation. The workload is mostly ORC scans and joins against dimension tables of a few hundred megabytes.
-- 4 GB containers: 27 per node, roughly one per core after headroom
SET hive.tez.container.size=4096;
SET hive.tez.java.opts=-Xmx3276m -XX:+UseG1GC -XX:+ExitOnOutOfMemoryError; -- 80%
SET tez.runtime.io.sort.mb=1300; -- ~40% of the 3276 MB heap
SET tez.runtime.unordered.output.buffer.size-mb=330; -- ~10% of the heap
SET hive.auto.convert.join.noconditionaltask.size=1145044992; -- 1092 MB, a third of the heap
-- One query joins a 2.5 GB dimension and keeps dying: give it bigger containers only
SET hive.tez.container.size=8192;
SET hive.tez.java.opts=-Xmx6553m -XX:+UseG1GC -XX:+ExitOnOutOfMemoryError;The arithmetic: 110,592 / 4,096 = 27 containers per node. A heap of 3,276 MB leaves about 820 MB for metaspace, threads and native buffers, which is comfortable for ORC readers with default codecs. The map-join threshold is in bytes. Together the sort buffer, unordered buffer and map-join budget come to about 2,720 MB, leaving roughly 550 MB of heap for operator state, readers and decompression buffers. The on-disk estimate behind the map-join budget understates the in-memory hash table, so if a task that sorts also loads a large map join, lower one of the two rather than trusting the arithmetic.
For the one heavy query, the session-level override doubles the container and heap, halving that query's parallelism on the node but leaving everyone else unchanged. If many queries need 8 GB, change the default instead and accept fewer containers per node, because a wave of retried OOMs costs more than a smaller cluster-wide parallelism.
The Tez Application Master
The Tez AM holds the DAG, the state of every task, task events and counters. Its heap grows with the number of tasks, so queries over tables with tens of thousands of small files or partitions generate many splits and can exhaust a default-sized AM. Symptoms are a query that fails before any task runs, or an AM container killed during scheduling. Raise tez.am.resource.memory.mb (and keep the -Xmx in tez.am.launch.cmd-opts at about 80% of it), but also look at the cause: compacting small files or grouping splits reduces the AM's load and speeds up the query. Container reuse (tez.am.container.reuse.enabled) means task JVMs live across tasks, so a slow leak in a UDF shows up as later tasks in the same container failing, not the first one; Hive on Tez execution covers that model.
HiveServer2 and Metastore heaps
HiveServer2's heap is set through the environment of its start script, usually a block in hive-env.sh keyed on the service name, or the equivalent field in Cloudera Manager or Ambari, which generate that file:
if [ "$SERVICE" = "hiveserver2" ]; then
export HADOOP_HEAPSIZE=16384 # MB
export HADOOP_OPTS="$HADOOP_OPTS -XX:+UseG1GC -XX:MaxGCPauseMillis=200 \
-XX:+HeapDumpOnOutOfMemoryError -XX:HeapDumpPath=/var/log/hive/hs2.hprof"
# GC log, JDK 9 and later (unified logging):
export HADOOP_OPTS="$HADOOP_OPTS \
-Xlog:gc*:file=/var/log/hive/hs2-gc.log:time,uptime,level,tags:filecount=10,filesize=50m"
# Java 8 instead, because it rejects -Xlog:
# export HADOOP_OPTS="$HADOOP_OPTS -XX:+PrintGCDetails -XX:+PrintGCDateStamps -Xloggc:/var/log/hive/hs2-gc.log"
fi
if [ "$SERVICE" = "metastore" ]; then
export HADOOP_HEAPSIZE=8192
fiHiveServer2 memory scales with concurrent sessions and with what those sessions compile and fetch. Planning a query over a table with hundreds of thousands of partitions builds large partition and statistics objects; a client that fetches millions of rows through JDBC can buffer large result batches; and every open session holds configuration and operation state. Cap hive.server2.thrift.max.worker.threads to bound concurrency, set hive.metastore.limit.partition.request so a careless SELECT without a partition filter fails fast instead of loading every partition, and close idle sessions with hive.server2.idle.session.timeout. The G1 garbage collector guide explains the pause behaviour behind the GC flags.
Long GC pauses in HiveServer2 have a nasty side effect when dynamic service discovery is on: HiveServer2 registers in ZooKeeper, and a pause longer than the ZooKeeper session timeout expires its registration, so clients stop being routed to a server that is, a few seconds later, perfectly healthy. If instances disappear from discovery under load, read the GC log before touching ZooKeeper.
The Metastore's heap is driven by partition-heavy calls and statistics. Repeated Metaspace exhaustion in HiveServer2, rather than heap exhaustion, usually points to class loading by session-added UDF JARs that are never unloaded; restart cycles hide it, and moving common UDFs into permanent functions or the auxiliary path fixes it.
Garbage collectors and flags by JDK
Which collector and which flags apply depends on the JDK, so check that first: Hive 3 and Hive 4.0 run on Java 8, and Hive 4.1, released in July 2025, added JDK 17 support. On Java 8 the default collector is Parallel, which maximizes throughput with long stop-the-world pauses; on JDK 9 and later the default is G1. For task containers, throughput matters more than pauses, so Parallel or G1 are both reasonable; for HiveServer2 and the Metastore, which serve interactive clients, G1 with a pause target is the usual choice.
GC logging flags changed in JDK 9. On Java 8 use -XX:+PrintGCDetails -XX:+PrintGCDateStamps -Xloggc:<file>; on JDK 9 and later use unified logging, -Xlog:gc*:file=<file>:time,uptime,level,tags. The two sets do not mix: Java 8 refuses to start with -Xlog, and on newer JDKs some old flags such as -XX:+PrintGCDateStamps were removed and also stop the JVM from starting, while others such as -Xloggc are only deprecated and remapped with a warning. Check the start-up log after every JDK upgrade. Always add -XX:+HeapDumpOnOutOfMemoryError to the services, with a dump path on a disk large enough to hold the heap, and consider -XX:+ExitOnOutOfMemoryError for task JVMs so a broken container dies instead of limping on.
Diagnosing kills and OutOfMemoryErrors
| Message or symptom | Where | Likely cause | Fix |
|---|---|---|---|
| Container is running beyond physical memory limits ... Killing container (exit code 143) | Tez task or AM | Heap plus off-heap exceeded the container; often -Xmx too close to container size | Lower -Xmx to about 80%, or raise the container |
java.lang.OutOfMemoryError: Java heap space in a map-join operator | Tez task | Small table larger in memory than estimated, or stale stats | Recompute statistics, lower the map-join threshold, or raise the container for that query |
GC overhead limit exceeded | Any | Heap nearly full, collector reclaiming almost nothing | Same as heap OOM; find what is retaining memory |
| Query fails before tasks start; AM killed | Tez AM | Too many splits or tasks | Raise AM memory; compact small files |
| HS2 pauses, clients time out, instance leaves ZooKeeper | HiveServer2 | Full or long GC under many sessions or large plans | GC log, heap dump, cap concurrency and partition requests |
OutOfMemoryError: Metaspace | HiveServer2 | Classes from session UDF JARs accumulating | Permanent functions, auxiliary JAR path |
Find the failing process before changing anything. The query's YARN application logs show task and AM failures; HiveServer2 and Metastore logs live on their hosts. Change one setting at a time, rerun the same query, and compare GC logs and Tez counters such as spilled records, not just wall-clock time. For Impala's very different memory model, see Impala memory limits.
What to do next
- Inventory every Hive JVM in your cluster with its current heap, container size and collector; note the JDK version of each.
- Check that
hive.tez.java.optseither omits-Xmx(letting the 0.8 fraction apply) or sets it to roughly 80% ofhive.tez.container.size. - Make container sizes multiples of
yarn.scheduler.minimum-allocation-mb. - Size the sort buffer, unordered buffer and map-join threshold from the container, then measure spills and failures on representative queries.
- Turn on GC logging with the right syntax for your JDK and heap dumps on OOM for HiveServer2 and the Metastore.
- Set
hive.metastore.limit.partition.request, a worker-thread cap and an idle-session timeout on HiveServer2. - For each OOM ticket, identify the process from the logs, match the message against the table above, and fix one setting at a time with a session-level override first.