Hive on Spark ran HiveQL queries as Spark jobs, with one remote Spark driver per HiveServer2 session. Hive 4.0 removed it, but many clusters still run Hive 2.x or 3.x with it, notably CDH 6 deployments, and they will until a migration is finished. Those clusters still need tuning, and the defaults are poor: a fresh installation runs few map joins, writes many small files and lets each session grab executors that others cannot use.
This article is a tuning guide for those clusters. How the engine compiles queries, its client timeouts and the migration path are covered in Hive on Spark Engine, in depth; this page assumes that background and works through a sizing example, session concurrency, parallelism, map joins, skew and small files, then gives a symptom-to-cause table. The worked numbers follow Cloudera's CDH 6 tuning guide, with the arithmetic redone.
Where the time goes in a Hive on Spark query
A Hive on Spark query spends its time in five places, and each has its own control:
- Session start. The first query in a session submits a Spark application to YARN and waits for a driver and executors. Controlled by session lifetime and pre-warming.
- Map stages. Read input splits. The number of tasks equals the number of splits from
CombineHiveInputFormat, which groups small underlying splits together. - Shuffle stages. Repartition data for joins and aggregations. The number of reduce-side tasks comes from Hive's reducer estimate.
- Joins. A shuffle join moves both sides; a map join broadcasts the small side and moves nothing. Controlled by statistics and the map-join threshold.
- Writes. Each final task writes at least one file per output partition. Controlled by merge settings and by how rows are distributed before the write.
Tune in that order of evidence, not of convenience: open the Spark UI for a slow query, find the stage that dominates, and change the control that owns that stage.
Sizing executors: a worked example
Start with YARN, because every Spark setting is carved out of it. Cloudera's example host has 32 cores and 120 GB. Reserve one core each for the NodeManager and DataNode and two for the operating system, leaving 28; reserve 20 GB for the same processes, leaving 100 GB. Set yarn.nodemanager.resource.cpu-vcores=28 and yarn.nodemanager.resource.memory-mb to 100 GB.
Then size executors. Cloudera recommends 4, 5 or 6 cores per executor, partly because the HDFS client can misbehave with many concurrent writers in one JVM. Pick the value that divides the YARN cores evenly: 28 divided by 4 is 7 executors per host with no idle cores, where 6 would waste 4 cores and 5 would waste 3. Divide memory the same way: 100 GB over 7 executors is about 14 GB each, split into heap and overhead.
-- per session or in hive-site.xml (Hive 2.x / 3.x)
set spark.executor.cores=4;
set spark.executor.memory=12g;
set spark.yarn.executor.memoryOverhead=2048; -- MB; spark.executor.memoryOverhead on newer Spark
set spark.driver.memory=10500m;
set spark.yarn.driver.memoryOverhead=1536; -- driver total about 12 GBCheck the arithmetic before applying it. Heap plus overhead, 14 GB, must be below yarn.scheduler.maximum-allocation-mb or YARN refuses the container. Four tasks share one heap, so a task has about 3.5 GB on average, and that is the budget map-join hash tables and aggregation buffers must fit in. Across 40 hosts the cluster holds 7 x 40 = 280 executors and 1,120 concurrent tasks. Cloudera's guide prints 160 at this step, which multiplies by the cores per executor rather than executors per host; work it out for your own hardware. For the driver, Cloudera's rule gives 12 GB in total when the NodeManager has more than 50 GB, with 10 to 15 percent as overhead.
Sessions, dynamic allocation and pre-warming
Every HiveServer2 session owns a Spark application, which changes capacity planning. Static allocation (spark.executor.instances) gives a session its executors for its whole life, idle or not, so twenty BI connections can pin the cluster. Cloudera recommends dynamic allocation for shared clusters, which needs the external shuffle service so executors can leave without losing shuffle files:
set spark.dynamicAllocation.enabled=true;
set spark.shuffle.service.enabled=true;
set spark.dynamicAllocation.minExecutors=0;
set spark.dynamicAllocation.maxExecutors=40; -- per session cap, derived below
set spark.dynamicAllocation.executorIdleTimeout=60s;Derive the cap from the queue, not from a guess. Suppose Hive gets a YARN queue of 60 percent of the cluster, about 168 executors, and at peak four heavy queries run at once. A cap of 40 per session lets four of them run at full speed with room left for light queries. Cloudera's guidance is that half the cluster's executors usually gives one query good performance, so a cap above half rarely pays. Remember the drivers too: twenty open sessions hold twenty 12 GB driver containers, 240 GB, whether or not they are running anything. Close idle sessions in BI tools and keep connection pools small.
Pre-warming helps short sessions such as scheduled jobs. Without it, a job starts before all executors are ready, and because Hive considers the executors available at submission when choosing reducer counts, the first query can run with too little parallelism. Set hive.prewarm.enabled=true and hive.prewarm.numcontainers (default 10) up to the session's executor cap. If the cluster cannot supply that many, the wait can last up to 30 seconds, so do not set it higher than the queue can provide.
Parallelism at the shuffle boundary
On the map side, parallelism follows the input splits; if a scan has too few tasks for its data, the files are probably large and unsplittable, or the split size settings group them too aggressively. The more common problem is on the shuffle side. Hive estimates the reduce-side task count from the data entering the stage divided by hive.exec.reducers.bytes.per.reducer, bounded by hive.exec.reducers.max (default 1009), and takes the available executors into account. An explicit mapreduce.job.reduces overrides the estimate.
Work an example. A join stage receives an estimated 180 GB. With the Apache default of 256 MB per reducer, Hive plans about 700 tasks; with 64 MB, the value in Cloudera's recommended list, it plans about 2,700 and hits the 1009 cap. A session capped at 40 executors has 160 task slots, so 700 tasks run in about four and a half waves, which is healthy. Cloudera's experience is that Spark is less sensitive to this setting than MapReduce was, provided there are enough tasks to keep every executor busy. Lower it when stages show fewer tasks than slots or when tasks spill heavily; raise it when thousands of tiny tasks spend more time scheduling than working.
-- Inspect the plan before and after a change
EXPLAIN
SELECT c.region, SUM(o.amount)
FROM orders o JOIN customers c ON o.customer_id = c.id
WHERE o.order_date >= '2026-09-01'
GROUP BY c.region;
-- In the Spark section look for: Reducer N <- Map M (PARTITION-LEVEL SORT, 703)
-- the number in brackets is the planned task count for that edge.
Map joins and the rawDataSize trap
Map joins are the largest single lever, and Hive on Spark makes them harder to get than MapReduce did. The threshold hive.auto.convert.join.noconditionaltask.size is compared with table statistics, and the two engines read different statistics. Hive on MapReduce uses totalSize, the size on disk. Hive on Spark uses rawDataSize, the estimated size in memory, when it is available. For compressed columnar tables the second is often many times the first.
So a dimension table that is 15 MB of ORC on disk and 150 MB raw becomes a map join under MapReduce with Cloudera Manager's 20 MB setting, and a shuffle join under Spark. The Apache default in HiveConf is lower still, 10 MB. Cloudera suggests raising the threshold to around 200 MB for Spark. Before you do, check memory: the small sides of a chain of map joins are held in memory as hash tables, which are larger than the raw data, and that space comes out of the roughly 3.5 GB a task can count on. Raise the threshold in steps and watch executor GC time and failures.
ANALYZE TABLE customers COMPUTE STATISTICS;
ANALYZE TABLE customers COMPUTE STATISTICS FOR COLUMNS;
DESCRIBE FORMATTED customers; -- read numRows, totalSize, rawDataSize
set hive.stats.fetch.column.stats=true;
set hive.auto.convert.join=true;
set hive.auto.convert.join.noconditionaltask=true;
set hive.auto.convert.join.noconditionaltask.size=200000000; -- bytes; Spark compares rawDataSizeMissing statistics are the most common reason a join that should be a map join is not; tables loaded by external tools often have none. Keep hive.stats.autogather on for inserts and run ANALYZE after bulk loads. Join strategy selection in general is covered in Hive join strategy selection.
Skewed keys
Skew shows up in the Spark UI as a stage where most tasks finish in seconds and a few run for many minutes, with shuffle read for those tasks far above the median. The cause is a key with a disproportionate share of rows, often a null, a default value such as 'UNKNOWN', or one very large customer. Find it directly:
SELECT customer_id, COUNT(*) AS n
FROM orders WHERE order_date >= '2026-09-01'
GROUP BY customer_id ORDER BY n DESC LIMIT 20;For aggregations, hive.groupby.skewindata=true adds a first stage that distributes rows randomly and partially aggregates them before the final aggregation by key; it costs an extra shuffle and does not apply to every query shape. For joins, the dependable fix is in the query: filter null keys out of the join if they cannot match anyway, and handle the few hot keys separately with a map join, combining the results with UNION ALL. Hive's runtime skew-join option was designed around MapReduce; test it on your build before relying on it. Hive skew join optimization goes further.
Small output files
Each final task writes its own file, and with dynamic partitioning it can write one per partition it sees. A 700-task insert into 30 date partitions can create 21,000 files, which loads the NameNode and slows every later scan. Merging is off by default on Spark (hive.merge.sparkfiles=false); Cloudera recommends turning it on:
set hive.merge.sparkfiles=true;
set hive.merge.smallfiles.avgsize=16000000; -- merge if average output file is below 16 MB
set hive.merge.size.per.task=256000000; -- target size of merged files
INSERT OVERWRITE TABLE orders_daily PARTITION (order_date)
SELECT ..., order_date FROM staging_orders
DISTRIBUTE BY order_date; -- each partition written by fewer tasksMerging adds a job at the end of the query; DISTRIBUTE BY the partition column avoids creating the files at all, at the risk of skew when one partition is much larger than the rest. Hive small files covers compaction of files that already exist.
Symptoms and likely causes
| Symptom | Likely cause | First thing to try |
|---|---|---|
| First query in a session is slow, later ones fast | Executor start-up | Pre-warm, or keep sessions open longer |
| Few long tasks in one stage | Key skew | Find hot keys; skewindata or split hot keys |
| Container killed for exceeding memory limits | Off-heap use above overhead | Raise memoryOverhead, keep the total under the YARN maximum |
| High GC time, executors lost in join stages | Map-join hash tables too large | Lower the threshold or fix statistics |
| Expected map join runs as shuffle join | No stats, or rawDataSize above the threshold | ANALYZE, then raise the threshold in steps |
| Fewer tasks than executor slots | Bytes per reducer too high, or few splits | Lower bytes per reducer; check input file sizes |
| Thousands of tiny output files | Merge off, many tasks per partition | Enable merge; DISTRIBUTE BY partition column |
| Session fails to start under load | Queue full, handshake timeout | See the timeout settings in the engine article |
Trade-offs
Every lever trades something. Larger executors allow bigger map joins but lengthen GC pauses and strand more memory when idle. Dynamic allocation shares the cluster but adds start-up latency each time executors are re-acquired. More reducers improve balance but increase scheduling overhead and output files. A higher map-join threshold removes shuffles but turns a statistics error into an executor failure. Tune with the Spark UI open and change one setting at a time; and remember that every hour spent here is spent on an engine Hive 4 no longer ships, so fix what hurts and put the rest of the effort into the migration.
What to do next
- Record your Hive version and confirm which jobs set
hive.execution.engine=spark. - Redo the node-carving arithmetic for your hardware: YARN cores and memory, executors per host, heap and overhead, per-task memory, driver size.
- Enable dynamic allocation with the external shuffle service, and set a per-session executor cap derived from the queue size and peak concurrency.
- Turn on pre-warming for scheduled jobs, capped at the session's executor limit.
- Run
ANALYZEon dimension tables, compare totalSize with rawDataSize, and raise the map-join threshold in steps while watching GC time. - Check EXPLAIN task counts for your ten heaviest queries against the session's task slots.
- List the top keys of every large join and group-by column, and split or filter the hot ones.
- Enable Spark file merging and add DISTRIBUTE BY to dynamic-partition inserts.
- Keep a before-and-after record of runtime and container-hours for every change, and reuse it to validate the migration to Tez or Spark SQL.