Hive on Spark was an execution engine for Apache Hive that ran HiveQL queries as Apache Spark jobs. You switched to it with one setting, hive.execution.engine=spark, and the same SQL, metastore, SerDes and UDFs kept working while the physical execution moved from MapReduce to Spark. For several years it was a supported engine in Hive 2.x and 3.x and a common choice in Hadoop distributions.

It is also gone. In Hive 4.0.0 the validator for hive.execution.engine accepts only mr and tez; Spark was removed from the main branch (tracked as HIVE-26134) because it was not being maintained. That makes this article two things: an explanation of how the engine actually worked, for the many clusters still running it, and a practical guide to operating it safely until you move off it, and to moving off it. Configuration defaults quoted here were read from the Hive 3.1.3 source.

Advertisement

Why an engine swap was possible at all

Hive separates what a query means from how it runs. The parser and semantic analyser build an operator tree: table scans, filters, joins, group-bys, reduce sinks and file sinks. The cost-based optimiser rewrites that tree. Only after that does a compiler for the chosen engine cut the tree into units of work. For MapReduce the units are map and reduce phases of separate jobs, with every intermediate result written to HDFS. For Tez, described in Hive on Tez execution, they are vertices in one DAG.

Hive on Spark reused the same idea. Its SparkCompiler produced a SparkWork: a graph of MapWork and ReduceWork nodes, the same work classes MapReduce used, connected by edges that describe the shuffle between them. Because the operator implementations were shared, most Hive features worked on Spark without being rewritten. What changed was how the work graph was executed and where the executors came from.

The architecture

Beeline / JDBC clientone HS2 sessionHiveServer2parse, CBO, SparkCompilerSparkWork: graph of MapWork / ReduceWorkSparkSession per HS2 sessionRemoteHiveSparkClientRPCRemote Spark driverYARN cluster modeExecutorsHiveMapFunction, HiveReduceFunctionshuffleSpark shuffleGROUP / SORT / PARTITION-LEVEL SORTHDFS / object storetable files, resultsspark-submit launches thedriver on first queryYARN RMcontainers
A HiveServer2 session owns a SparkSession object that talks over RPC to a remote Spark driver running as its own YARN application. The driver runs Hive's map and reduce functions on executors, connected by Spark shuffles.

The most important fact about Hive on Spark is in the middle of this diagram. When spark.master is not local, each HiveServer2 session gets its own SparkSession, implemented with RemoteHiveSparkClient. On the first query in the session, Hive launches a separate Spark application through spark-submit. Its driver, the RemoteDriver, connects back to HiveServer2 over an RPC channel, and all later queries in that session are submitted to the same driver as Spark jobs. The factory defaults in Hive 3.1.3 are spark.master=yarn, deploy mode cluster, the application name "Hive on Spark" and the Kryo serializer.

In local mode, by contrast, all sessions share one in-process Spark context, which is useful for tests and nothing else. The per-session design gives isolation, since one session's driver crash does not affect another, and it amortises Spark start-up across all the queries in a session. It also means a long-lived session holds a YARN application, and possibly executors, even when it is idle. That single property explains most of the operational behaviour below.

Advertisement

How work becomes RDDs

Inside the driver, Hive's SparkPlanGenerator turns the SparkWork into Spark transformations. Each MapWork reads its input as a Hadoop RDD, using the table's input format, and runs HiveMapFunction over each partition. That function drives the same operator pipeline a MapReduce mapper would: deserialize rows, filter, project, probe map-join hash tables, pre-aggregate, and emit keyed rows through the reduce sink. Each ReduceWork runs HiveReduceFunction over the grouped or sorted output of a shuffle.

The edge type decides which Spark shuffle is used. EXPLAIN prints it as the shuffle type:

Shuffle typeSpark operationUsed for
GROUPgroupByKey (when hive.spark.use.groupby.shuffle is true, the default)Group-bys where order within a key does not matter.
PARTITION-LEVEL SORTrepartitionAndSortWithinPartitionsJoins and group-bys that expect MapReduce-style sorted input per reducer; also GROUP when the setting above is false.
SORTsortByKeyGlobal ORDER BY, usually with one partition.

The Hive 3.1.3 description of hive.spark.use.groupby.shuffle states the trade-off plainly: groupByKey performs better but uses unbounded memory, so turn the setting off when there is a memory issue. Unlike MapReduce, the map and reduce stages of one SparkWork run as a single Spark job, and intermediate data stays in Spark's shuffle files rather than being written to HDFS between stages. For how those shuffle files are written and fetched, see Spark shuffle.

Worked example: reading a plan

EXPLAIN
SELECT c.region, sum(o.amount) AS revenue
FROM   orders o JOIN customers c ON o.customer_id = c.id
WHERE  o.dt = '2026-09-29'
GROUP  BY c.region
ORDER  BY revenue DESC;

-- abbreviated plan shape (Stage-1 depends on Stage-2)
Stage: Stage-1
  Spark
    Edges:
      Reducer 2 <- Map 1 (GROUP, 24)
      Reducer 3 <- Reducer 2 (SORT, 1)
    Vertices:
      Map 1       -- scan orders partition, map join with broadcast customers, partial aggregate
      Reducer 2   -- final aggregate per region
      Reducer 3   -- global order by, one partition
Stage: Stage-2
  Spark
    Vertices:
      Map 4       -- scan customers, write its hash table to the file system

Walk through the plan. The small customers table fits under the map-join threshold, so Hive runs a separate small Spark job that scans it and writes its hash table to the file system, and the executors of the main job load it; there is no shuffle for the join. Map 1 scans only the partition dt='2026-09-29' because of static partition pruning, probes the hash table, and partially aggregates by region. The GROUP edge with 24 partitions sends rows by region to Reducer 2, which finishes the sums. Reducer 3 sorts globally with a single partition, which is fine for a few dozen regions and a bottleneck for millions of rows.

Three things to check in any Hive on Spark plan: whether joins you expect to be map joins really are, since a missed map join adds a full shuffle of the large table; the partition count on each edge, which Hive estimates from statistics (hive.spark.use.op.stats defaults to true) so stale statistics give bad parallelism; and any SORT edge with one partition over a large input.

Configuration that matters

Hive on Spark needs a Spark build that does not include Hive's own jars, and a Spark version that matches the Hive release; the project's compatibility table pairs Hive 3.0.x with Spark 2.3.0 and Hive 2.3.x with Spark 2.0.0. Mismatches show up as class-loading errors when the remote driver starts. A typical session-level setup looks like this:

-- hive-site.xml or per session (Hive 2.x / 3.x only)
set hive.execution.engine=spark;
set spark.master=yarn;                     -- the HiveSparkClientFactory default
set spark.submit.deployMode=cluster;       -- also the default: driver runs in a YARN container
set spark.executor.cores=4;
set spark.executor.memory=12g;
set spark.executor.memoryOverhead=2g;      -- older Spark: spark.yarn.executor.memoryOverhead
set spark.dynamicAllocation.enabled=true;
set spark.shuffle.service.enabled=true;    -- required by dynamic allocation on Spark 2.x
set spark.dynamicAllocation.maxExecutors=40;

The Hive-side settings that control the client and driver, with their Hive 3.1.3 defaults:

PropertyDefaultWhat it controls
hive.spark.client.connect.timeout1000msTime for the remote driver to connect back to HiveServer2.
hive.spark.client.server.connect.timeout90000msHandshake timeout between HiveServer2 and the driver, checked on both sides; this must cover YARN queueing and driver start-up.
hive.spark.client.future.timeout60sTimeout for requests from Hive to the remote driver.
hive.spark.job.monitor.timeout60sTimeout for the job monitor to obtain a job's state.
hive.spark.client.rpc.max.size50 MBLargest RPC message between Hive and the driver.
hive.spark.dynamic.partition.pruningfalseDynamic partition pruning for joins on partition keys.
hive.spark.dynamic.partition.pruning.max.data.size100 MBCap on the data collected for dynamic pruning.
hive.spark.job.max.tasks-1Cancel jobs with more tasks than this; -1 means no limit.
hive.merge.sparkfilesfalseMerge small output files at the end of a job.

Dynamic partition pruning being off by default surprises people coming from Tez; queries that join a fact table to a filtered dimension on the partition key scan every partition unless you enable it. Small output files are the other common surprise: with many reducers and hive.merge.sparkfiles off, inserts can create thousands of tiny files.

Operating it: sessions, executors and queues

Because each HiveServer2 session owns a YARN application, capacity planning is per session. With static executors, twenty idle BI connections can pin twenty applications' worth of memory. Enable dynamic allocation with the external shuffle service so idle sessions release executors, set a maximum per session, and put Hive on Spark applications in their own YARN queue so they cannot starve other workloads. Connection pools in BI tools deserve attention: a pool of fifty connections means up to fifty drivers.

The first query in a session pays for spark-submit, YARN scheduling and driver start-up, which is typically seconds and can be much longer on a busy cluster. Users experience this as "the first query is slow". Pre-warming by running a trivial query at session start moves the cost but does not remove it. When YARN is busy the handshake can exceed hive.spark.client.server.connect.timeout, and the query fails with an error about creating the Spark client; raise the timeout, and fix the queue capacity that caused the wait.

Monitor at two levels. In HiveServer2 (see HiveServer2), track open sessions, query latency and client creation failures. In the Spark UI of each driver, track stage skew, shuffle spill and executor loss. Keep driver logs; failures inside the remote driver appear there, not in HiveServer2 logs.

Failure modes

  • Client creation timeouts. The driver did not start or connect back in time, usually because of YARN queue limits, a wrong Spark home, jar conflicts, or network rules blocking the driver's connection to HiveServer2.
  • Executor out-of-memory in GROUP shuffles. groupByKey holds all values for a key in memory; a skewed key kills executors. Turn off hive.spark.use.groupby.shuffle for the job and address the skew.
  • Map-join hash tables too large. A statistics error can broadcast a table that does not fit, failing executors; keep statistics fresh and set the map-join size threshold to what executor memory can hold.
  • Driver death takes the session with it. Every later query in the session fails until the client reconnects and a new driver starts.
  • Leaked applications. Sessions not closed cleanly can leave driver applications running; audit YARN for long-lived "Hive on Spark" applications.
  • Upgrade breakage. On Hive 4 the setting fails validation, so any script that pins the engine to Spark stops working.

Moving off Hive on Spark

There are two destinations. The first is Hive on Tez, the supported engine in Hive 3 and 4. It keeps HiveQL, ACID tables, HiveServer2, Ranger policies and UDFs unchanged, and adds features such as LLAP described in Hive LLAP. The second is Spark SQL reading the Hive metastore. Spark SQL is a separate engine with its own optimiser and SQL dialect; it does not run Hive's operators, and full Hive ACID table support is not built into Spark, so check each table type and each UDF before choosing it.

Either way, treat it as a query migration, not a configuration change:

# 1. Find every script and job that pins the engine
grep -rniE "hive\.execution\.engine\s*=\s*spark" etl/ oozie/ airflow/ conf/

# 2. Run the same query on both engines against a copy of the inputs
beeline -u "$JDBC" --hiveconf hive.execution.engine=spark -f q.sql > spark.out
beeline -u "$JDBC" --hiveconf hive.execution.engine=tez   -f q.sql > tez.out

# 3. Compare results order-insensitively, then compare runtime and container-hours
sort spark.out | sha256sum; sort tez.out | sha256sum

Compare results, not just success, because engines can differ in edge cases such as implicit casts or the order of results without ORDER BY. Retune after switching: Tez container sizes and reducer counts are set differently from Spark executors, and a query tuned for one engine can be slow on the other. Migrate in waves ordered by business criticality, and remove Spark-specific settings from hive-site.xml last. The wider list of removed features in Hive deprecated features is worth reviewing in the same upgrade.

Trade-offs, in hindsight

Hive on Spark delivered real gains over MapReduce: in-memory shuffles between stages, one job per query, and reuse of a warm driver within a session. Its weaknesses were structural. Every session carried a separate application, it depended on matching Spark versions, and it duplicated effort with Spark SQL, which grew its own Hive compatibility. When the maintainers had to choose, Tez stayed and Spark went.

What to do next

  1. Find out which Hive version you run and whether any job sets hive.execution.engine=spark; the grep above finds most of them.
  2. If you still run it, enable dynamic allocation, cap executors per session, and give it a dedicated YARN queue.
  3. Check the four settings that bite most often: client handshake timeout, dynamic partition pruning, groupby shuffle, and small-file merging.
  4. Run EXPLAIN on your ten most expensive queries and verify map joins, edge partition counts and single-partition sorts.
  5. Choose a target engine, Tez or Spark SQL, per workload, and run both engines side by side on copies of the inputs.
  6. Plan the Hive 4 upgrade only after no job depends on the Spark engine.
Key takeaway: Hive on Spark compiled HiveQL into a SparkWork graph of map and reduce work, ran it through a remote Spark driver owned by each HiveServer2 session, and connected stages with GROUP, PARTITION-LEVEL SORT or SORT shuffles. Its per-session applications, version coupling and memory-hungry groupByKey shuffles drove most operational pain, and Hive 4.0 removed it. If you still run it, bound its resources and check its defaults; then migrate to Tez or Spark SQL with side-by-side result comparison.