Hive LLAP (Live Long and Process) exists for one reason: to make Hive answer interactive SQL in seconds instead of minutes. It does that with long-running daemons that keep executors warm and cache column data off-heap, so a query does not pay for container start-up, JVM warm-up or a full re-read of ORC files. When an LLAP cluster is slow, it is almost never because LLAP is slow. It is because the daemons are sized badly, the cache is not being hit, queries queue behind each other, or the benchmark was measuring something other than what people thought.
This page is about performance work, not architecture. The components themselves, the IO elevator and the end-to-end path of a query are explained in Hive LLAP architecture. Here we start from what a fast query needs, work through a concrete memory budget for a 256 GB node, choose executor and IO-thread counts, make the cache effective, read the counters that prove it, build an honest benchmark and use workload management to protect latency. Configuration names are from Apache Hive 3 and 4; sizing figures are rules of thumb to validate on your own data, not documented constants.
Where interactive query time goes
Break one query's wall time into stages before tuning anything. HiveServer2 compiles the SQL and runs the cost-based optimizer, which needs metastore calls and table statistics. A Tez application master accepts the DAG; HiveServer2 can keep a pool of pre-started AMs so this costs milliseconds rather than the several seconds of a YARN application launch. The AM hands fragments to the LLAP task scheduler, which places each on a daemon. Inside the daemon a fragment may wait in a queue for a free executor. Then it reads data, either from the off-heap cache or from HDFS or object storage, executes vectorized operators and shuffles output to the next vertex.
Each stage has a different fix. Slow compilation means missing statistics, huge partition counts or a slow metastore. AM wait means the Tez session pool is too small for the concurrency. Queue wait means too few executors or one heavy query hogging them. Slow reads mean cache misses. Slow execution means too little parallelism, spilling hash joins or plans without vectorization. Measure before tuning.
The daemon's memory, in one picture
Every LLAP daemon is one YARN container on one node, normally one per node. Its memory is split three ways. The JVM heap (--xmx) holds the executors' working memory: operator state, hash tables for map joins, sort buffers and vectorized batches. The off-heap cache (--cache) holds decoded column chunks and ORC metadata in direct memory managed by a buddy allocator, outside garbage collection. Headroom is everything else the process needs: thread stacks, NIO buffers, JIT code and metaspace. The container size (--size) must cover all three, or YARN will kill the daemon for exceeding its limit.
Worked example: sizing a 256 GB node
Take ten worker nodes, each with 256 GB of RAM and 32 cores. Reserve memory for the operating system, the DataNode and the NodeManager; 60 GB is generous and leaves 196 GB for the LLAP container. Next decide concurrency. A widely used rule of thumb from Hortonworks and Cloudera sizing guidance is about 4 GB of heap per executor, enough for typical map-join hash tables. Sixteen executors on 32 cores leaves cores for IO threads and the shuffle handler, so heap is 16 x 4 = 64 GB.
Headroom is commonly taken as roughly 6 percent of heap, capped at a few gigabytes: about 4 GB here. The cache gets the rest: 196 - 64 - 4 = 128 GB per node, 1.28 TB across the cluster. Compare that with your hot working set, the columns and partitions dashboards actually read. A 600 GB hot set fits; a 5 TB one never will, and the realistic aim is caching metadata and the newest partitions.
def size_llap(node_ram_gb, yarn_reserved_gb, executors, gb_per_executor=4.0,
headroom_frac=0.06, headroom_cap_gb=6.0):
"""Rule-of-thumb LLAP daemon sizing. Validate every output on your own workload."""
container = node_ram_gb - yarn_reserved_gb # what YARN may give the daemon
xmx = executors * gb_per_executor # heap scales with concurrency
headroom = min(xmx * headroom_frac, headroom_cap_gb) # native memory outside heap and cache
cache = container - xmx - headroom
if cache < 0.2 * container:
raise ValueError(f"cache only {cache:.0f} GB: fewer executors or a bigger node")
return dict(size=f"{container:.0f}g", xmx=f"{xmx:.0f}g",
cache=f"{cache:.0f}g", executors=executors, iothreads=executors)
print(size_llap(node_ram_gb=256, yarn_reserved_gb=60, executors=16))
# {'size': '196g', 'xmx': '64g', 'cache': '128g', 'executors': 16, 'iothreads': 16}The helper encodes these rules and refuses budgets where the cache would be starved. Use it to compare options, then start the daemons with the result:
# Package and start the LLAP YARN application with the computed numbers
hive --service llap --name llap0 --instances 10 \
--size 196g --xmx 64g --cache 128g --executors 16 --iothreads 16
# Check it: running instances, per-daemon state, recent failures
hive --service llapstatus --name llap0
# hive-site.xml / session settings that decide whether queries use it well
set hive.execution.engine=tez;
set hive.execution.mode=llap;
set hive.llap.execution.mode=all; -- run every fragment in the daemons
set hive.llap.io.enabled=true;
set hive.llap.client.consistent.splits=true; -- same split goes to the same daemon
set hive.tez.exec.print.summary=true; -- prints DAG timings and LLAP IO counters
Executors, IO threads and the wait queue
Executors are LLAP's slots: each runs one fragment at a time. More executors means more concurrent fragments, but each needs heap and a core. Going past the core count rarely helps, because vectorized scans and aggregations are CPU-bound once data is cached. Too few executors and fragments wait. The symptom is fragments sitting in the daemon's wait queue while CPU is idle on other nodes, which usually means skewed placement rather than too few slots overall.
IO threads do the reading and decoding into the cache on behalf of executors; setting --iothreads equal to executors is the usual starting point. When fragments from several queries compete, the daemon orders its wait queue and can preempt work that is not guaranteed in favour of fragments that are, which keeps short queries from starving behind a long one. That is the right behaviour for interactive work. It also means a big ETL query on LLAP makes interactive latency less predictable, which is one argument for running heavy batch work in plain Tez containers (hive.llap.execution.mode=none for those sessions) or for confining it to a pool, as shown later.
Making the cache actually hit
A large cache is useless if the same data lands on a different daemon each time. With hive.llap.client.consistent.splits=true splits are mapped to daemons by consistent hashing, so a file region is read and cached by the same node every time, and adding or losing a node remaps only a fraction of splits. If the owning daemon stays busy, the scheduler eventually runs the split elsewhere as a cache miss, so a saturated cluster hits less than its cache size suggests.
Data layout matters as much as cache size. The cache stores column chunks from ORC stripes, so it rewards queries that read few columns, ORC files written with useful sort order so that min/max indexes and bloom filters skip stripes, and reasonably sized files rather than thousands of tiny ones. Run compaction on ACID tables so readers do not merge delta files on every query.
Eviction uses LRFU, a blend of least-recently-used and least-frequently-used; hive.llap.io.use.lrfu enables it. A single full scan of a cold historical table can push the dashboard's hot chunks out, which is why heavy scans are best kept off the interactive daemons or killed by a trigger. When RAM is short, Hive can back the cache with memory-mapped files on local SSD (hive.llap.io.allocator.mmap plus a path), giving a larger, slower tier that is still far faster than remote object storage.
Reading the evidence: counters and the daemon UI
With hive.tez.exec.print.summary=true beeline prints, after each query, a per-vertex time breakdown and an LLAP IO summary. The counters to read first are CACHE_HIT_BYTES and CACHE_MISS_BYTES (the hit ratio for this query), METADATA_CACHE_HIT (whether ORC footers came from memory) and SELECTED_ROWGROUPS (how many row groups survived predicate pushdown). A dashboard query with a high miss ratio on its second run tells you the cache or placement is not working. A query that selects nearly every row group tells you the data layout is not helping.
Each daemon serves a web UI, on port 15002 by default, with JMX metrics, configuration and thread stacks. Its JMX output exposes executor usage, wait-queue length and cache allocation, which you can scrape into your monitoring system. hive --service llapstatus shows whether every instance is running. A restarting daemon loses its cache, which shows up as latency spikes, not errors. See Hive monitoring for wiring these into dashboards and alerts.
An honest benchmark
LLAP benchmarks go wrong in predictable ways. One run after start-up measures cold reads, twenty in a row measure a perfectly hot cache, and one query at a time hides queueing. A useful benchmark records three numbers per query: the cold time (first run after a daemon restart), the hot median, and the 95th percentile under realistic concurrency.
#!/usr/bin/env bash
# Run each query N times at a fixed concurrency and record wall time per run.
# Run 1 after a daemon restart is "cold"; later runs are "hot".
JDBC="jdbc:hive2://hs2.example.com:10000/default;transportMode=binary"
for q in queries/*.sql; do
for run in 1 2 3 4 5; do
start=$(date +%s.%N)
beeline -u "$JDBC" --silent=true -f "$q" > /dev/null 2> "logs/$(basename $q).$run.err"
end=$(date +%s.%N)
echo "$(basename $q),$run,$(echo "$end - $start" | bc)" >> results.csv
done
doneRun the same script with several clients in parallel to get the concurrency figure, and keep the stderr logs, which contain the counters. Change one variable per run. Refresh statistics before every run. Hive's cost-based optimizer explains what the planner does with them, and vectorization covers checking with EXPLAIN that operators run in vectorized mode.
Workload management as a performance tool
On a shared LLAP cluster the biggest latency risk is another tenant. Hive 3 workload management divides executor capacity into pools with guaranteed fractions and concurrency limits, maps users or applications to pools, and runs triggers that move or kill queries crossing a threshold. Used well, it is a latency guarantee for dashboards:
CREATE RESOURCE PLAN daytime;
CREATE POOL daytime.bi WITH ALLOC_FRACTION = 0.7, QUERY_PARALLELISM = 20, SCHEDULING_POLICY = 'fair';
CREATE POOL daytime.etl WITH ALLOC_FRACTION = 0.3, QUERY_PARALLELISM = 4;
-- A dashboard query that runs long is not a dashboard query: move it out of the way
CREATE TRIGGER daytime.slow_bi WHEN ELAPSED_TIME > 30000 DO MOVE TO etl;
ALTER POOL daytime.bi ADD TRIGGER slow_bi;
-- Kill scans that would evict the whole cache
CREATE TRIGGER daytime.huge_scan WHEN HDFS_BYTES_READ > 500000000000 DO KILL;
ALTER POOL daytime.etl ADD TRIGGER huge_scan;
CREATE USER MAPPING 'tableau_svc' IN daytime TO bi;
ALTER RESOURCE PLAN daytime SET DEFAULT POOL = etl;
DROP POOL daytime.default; -- created with the plan at fraction 1.0; fractions must sum to at most 1
ALTER RESOURCE PLAN daytime VALIDATE;
ALTER RESOURCE PLAN daytime ENABLE ACTIVATE;Size the interactive pool for peak concurrency, push long runners to a batch pool and kill scans big enough to flush the cache. Compare Impala admission control.
Failure modes
| Symptom | Likely cause | What to do |
|---|---|---|
| Daemons killed by YARN, cache rebuilt repeatedly | Container smaller than heap + cache + headroom | Recompute the budget; raise headroom or shrink cache |
| Long full GC pauses, then slow queries | Heap per executor too small for map-join tables | Fewer executors or a lower map-join threshold |
| Second run as slow as the first | Consistent splits off, or busy daemons forcing remote runs | Enable consistent splits; check locality and wait queue |
| Latency spikes when ETL runs | Big scans evicting hot chunks and taking executors | Separate pool with kill trigger, or run ETL in plain Tez |
| Fast execution, slow total time | Compilation, metastore or Tez AM wait | Fix stats, partition counts and Tez session pool size |
| Rows read far above rows returned | No stripe skipping: unsorted data, missing bloom filters | Rewrite tables sorted by filter columns |
Trade-offs
- Cache versus concurrency: every gigabyte of heap given to executors is taken from the cache. Size for the concurrency you need, not the most you can fit.
- Always-on cost: LLAP holds memory and cores on every node all day. For bursty or batch-only use, autoscaled plain Tez or a managed warehouse service may cost less.
- Predictability versus utilization: workload management with strict pools leaves capacity idle at times, and that idle capacity is what buys predictable dashboard latency.
What to do next
- Turn on hive.tez.exec.print.summary and collect counters for your top twenty dashboard queries.
- Break their wall time into compile, AM wait, queue, read and execute, and fix the largest stage first.
- Recompute the daemon budget for your nodes with the sizing helper and compare the cache with the measured hot working set.
- Confirm consistent splits are enabled and check the hit ratio on second runs.
- Rewrite the hottest tables sorted on their main filter columns, with compaction running.
- Build the cold, hot and concurrent benchmark, and change one variable per run.
- Create a resource plan that protects interactive pools and kills cache-flushing scans.