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.

Advertisement

The JVMs in a Hive deployment

Start by knowing which process failed, because the same words, OutOfMemoryError, mean different things in each:

Beeline / JDBC clientsmall heap, large fetchesHiveServer2 JVMsessions, compile, fetchMetastore JVMpartitions, stats, locksTez AM JVMone per session/DAGSQLThriftsubmit DAGYARN NodeManager: container of hive.tez.container.size MBJava heap (-Xmx, ~80%)off-heapsort buffermap-join hash tableoperatorsmetaspace, stacks, NIOlaunch tasksYARN kills the process when total physical memory, heap plus off-heap, exceeds the container.The JVM throws OutOfMemoryError when the heap alone is exhausted. Different limits, different fixes.
The JVMs in a Hive-on-Tez deployment. Long-running services (HiveServer2, Metastore) are sized by their own heap settings; Tez task JVMs live inside YARN containers, where the heap is only part of what the container must hold.
JVMLifetimeMemory driven byMain settings
HiveServer2Long-running serviceConcurrent sessions, query compilation (plans over many partitions), result fetch buffers, operation logsHeap in hive-env.sh or the cluster manager; GC flags
MetastoreLong-running servicePartition listings, column statistics, object caches, concurrent Thrift callsHeap, GC flags, backing database pool
Tez AMPer session or per queryNumber of tasks and splits, DAG size, task events and counterstez.am.resource.memory.mb, tez.am.launch.cmd-opts
Tez taskPer container; reused across tasksSort and shuffle buffers, map-join hash tables, operator state, ORC/Parquet readershive.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.

Advertisement

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.size caps 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
fi

HiveServer2 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 symptomWhereLikely causeFix
Container is running beyond physical memory limits ... Killing container (exit code 143)Tez task or AMHeap plus off-heap exceeded the container; often -Xmx too close to container sizeLower -Xmx to about 80%, or raise the container
java.lang.OutOfMemoryError: Java heap space in a map-join operatorTez taskSmall table larger in memory than estimated, or stale statsRecompute statistics, lower the map-join threshold, or raise the container for that query
GC overhead limit exceededAnyHeap nearly full, collector reclaiming almost nothingSame as heap OOM; find what is retaining memory
Query fails before tasks start; AM killedTez AMToo many splits or tasksRaise AM memory; compact small files
HS2 pauses, clients time out, instance leaves ZooKeeperHiveServer2Full or long GC under many sessions or large plansGC log, heap dump, cap concurrency and partition requests
OutOfMemoryError: MetaspaceHiveServer2Classes from session UDF JARs accumulatingPermanent 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

  1. Inventory every Hive JVM in your cluster with its current heap, container size and collector; note the JDK version of each.
  2. Check that hive.tez.java.opts either omits -Xmx (letting the 0.8 fraction apply) or sets it to roughly 80% of hive.tez.container.size.
  3. Make container sizes multiples of yarn.scheduler.minimum-allocation-mb.
  4. Size the sort buffer, unordered buffer and map-join threshold from the container, then measure spills and failures on representative queries.
  5. Turn on GC logging with the right syntax for your JDK and heap dumps on OOM for HiveServer2 and the Metastore.
  6. Set hive.metastore.limit.partition.request, a worker-thread cap and an idle-session timeout on HiveServer2.
  7. 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.
Key takeaway: Hive runs several kinds of JVM, and each fails differently: Tez tasks inside YARN containers, the Tez AM, HiveServer2 and the Metastore. Keep task heaps at about 80% of the container so off-heap memory fits, size sort buffers and map-join thresholds from the container, and fix heavy queries with session-level overrides. Bound HiveServer2 with worker, partition and session limits, log GC with the right flags for your JDK, and always identify the failing process before changing a setting.