Every engine in a Hadoop-style lakehouse asks the Hive Metastore (HMS) the same questions: where does this table live, which partitions exist, what is the schema, which transactions are open. HiveServer2, Spark, Impala, Trino and Flink all call it before they read a byte. When it is down, every query fails at planning, every ingestion job fails at commit, and the dashboards go red together.
Making HMS highly available is less about HMS than people expect. The servers themselves are stateless Thrift services and easy to run several of. The hard parts are the database underneath, the background housekeeping threads that must run on exactly one server, and client behaviour during failover. This article covers each, then works through a three-server design and the failure drills that prove it. For the general HMS architecture, see the Hive Metastore overview.
Where the state actually lives
An HMS server is a Java process exposing a Thrift API, by default on port 9083, that translates calls such as get_table and add_partitions into SQL against a relational database through DataNucleus. Table definitions, partitions, column statistics, privileges, the notification event log, and the ACID transaction and lock tables all live in that database. The server keeps connection pools and caches, but nothing that cannot be rebuilt.
That gives the HA design its shape: run several HMS servers active-active, any of which can serve any request, in front of one database that must be highly available in its own right. Adding HMS servers adds throughput and tolerates server loss; it does nothing for database loss. A three-node HMS tier on one unreplicated MySQL instance has a single point of failure, just a less obvious one.
Client-side failover
HA starts in the client configuration. hive.metastore.uris (metastore.thrift.uris in the standalone metastore naming) takes a comma-separated list. With hive.metastore.uri.selection at its default of RANDOM, each client shuffles the list and connects to the first reachable server, which spreads load across the fleet without a load balancer. SEQUENTIAL always tries the list in order, which is useful for keeping clients on a primary site and using the rest as fallbacks.
If a call fails with a transport error, the client reconnects to another server and retries. The relevant settings, with their defaults on current Apache Hive, are hive.metastore.connect.retries (3 attempts to open a connection), hive.metastore.client.connect.retry.delay (1 second between attempts), hive.metastore.failure.retries (1 retry of a failed call) and hive.metastore.client.socket.timeout (600 seconds).
<!-- hive-site.xml / metastore-site.xml on every CLIENT: HiveServer2, Spark, Impala, Trino -->
<property>
<name>hive.metastore.uris</name>
<value>thrift://hms1.example.com:9083,thrift://hms2.example.com:9083,thrift://hms3.example.com:9083</value>
</property>
<property>
<name>hive.metastore.uri.selection</name>
<value>RANDOM</value> <!-- or SEQUENTIAL: always try the first URI first -->
</property>
<property>
<name>hive.metastore.connect.retries</name>
<value>3</value> <!-- attempts to open a connection -->
</property>
<property>
<name>hive.metastore.client.connect.retry.delay</name>
<value>1s</value>
</property>
<property>
<name>hive.metastore.failure.retries</name>
<value>1</value> <!-- retries of a failed Thrift call after reconnecting -->
</property>
<property>
<name>hive.metastore.client.socket.timeout</name>
<value>600s</value> <!-- lower it so a hung server is abandoned sooner -->
</property>
<property>
<name>hive.metastore.client.socket.lifetime</name>
<value>1800s</value> <!-- 0 = forever; a finite value rebalances long-lived clients -->
</property>Two defaults deserve attention. A 600-second socket timeout means a server that accepts connections but hangs, for example while blocked on a dead database connection, holds a query for ten minutes before failover. Many operators lower it, but not below the duration of your slowest legitimate call; large add_partitions or get_partitions on a table with a million partitions can take minutes. And hive.metastore.client.socket.lifetime defaults to 0, meaning a connection lives forever. Long-running clients such as HiveServer2 and Impala's catalog daemon therefore stay pinned to whichever server they picked at startup; after a rolling restart, the last server restarted gets almost no traffic. A finite lifetime makes clients reconnect periodically and rebalance.
Retries have a correctness edge. If a create_table or add_partition reached the database but the response was lost, the retry can fail with an "already exists" error for a change that actually succeeded. Ingestion code should treat that error on retry as possible success and verify, rather than failing the job.
Load balancers and ZooKeeper discovery
There are three ways for clients to find servers. The static URI list above is the simplest and needs no extra infrastructure, but every client's configuration must change when servers are added or renamed. A TCP load balancer in front of the servers gives clients one stable name. It has two catches: Thrift connections are long-lived, so the balancer balances connections, not requests; and with Kerberos, clients build the service principal from the hostname they connect to, so the load balancer's name must be covered by the HMS keytab, or authentication fails in a way that looks like a network error.
Current Apache Hive also supports ZooKeeper discovery. With metastore.service.discovery.mode set to zookeeper, the URIs name the ZooKeeper ensemble instead, HMS servers register themselves under a namespace (default hive_metastore), and clients pick from the live registrations. Servers can be added without touching clients. Check that every client engine and version you run supports it before relying on it; a static list is the lowest common denominator.
<!-- ZooKeeper-based discovery: the URIs now name the ZooKeeper ensemble, not HMS hosts -->
<property>
<name>metastore.service.discovery.mode</name>
<value>zookeeper</value>
</property>
<property>
<name>metastore.thrift.uris</name>
<value>zk1.example.com:2181,zk2.example.com:2181,zk3.example.com:2181</value>
</property>
<property>
<name>metastore.zookeeper.namespace</name>
<value>hive_metastore</value>
</property>
Housekeeping: the part that must not be active-active
HMS also runs background work: the ACID compaction initiator and cleaner, automatic partition discovery and retention, expiry of the notification event log, cleanup of the change-management directory used by replication, and aborted-transaction cleanup. These must run on one server. Two compaction initiators can queue duplicate requests and contend on the same tables; zero initiators means delta files pile up until reads slow to a crawl. See Hive compaction in depth for what the initiator decides.
Older deployments handled this by hand: set hive.compactor.initiator.on and hive.compactor.cleaner.on to true on exactly one server and false elsewhere. The flaw is that when that server dies, compaction stops until someone notices.
Leader election replaces that. metastore.housekeeping.leader.election takes host or lock. With host, the server whose name matches metastore.housekeeping.leader.hostname runs housekeeping; it is simple, but it does not fail over. With lock, servers compete for an exclusive Hive lock in the metastore database; the holder is the leader, and if it dies another server acquires the lock and takes over. The default differs between releases and distributions (Apache Hive master defaults to lock, while some vendor docs describe host-based election as the default), so set it explicitly.
<!-- metastore-site.xml on every HMS SERVER -->
<property>
<name>metastore.housekeeping.leader.election</name>
<value>lock</value> <!-- lock: elect via a Hive lock and fail over; host: fixed hostname -->
</property>
<!-- Only with election = host: must equal metastore.thrift.bind.host on the leader -->
<property>
<name>metastore.housekeeping.leader.hostname</name>
<value>hms1.example.com</value>
</property>
The database is the HA problem
Pick a database topology that loses no committed transactions on failover. Managed multi-AZ offerings, PostgreSQL with synchronous replication, or MySQL with semi-synchronous replication all qualify. Asynchronous replication with manual promotion does not: the last seconds of commits are lost, and those seconds are the metastore's view of which partitions and ACID transactions exist. After such a failover, a partition an ingestion job saw committed can vanish while its files remain, open transactions can be forgotten, and notification event IDs can go backwards, which breaks any consumer that tracks the last event it processed. That includes replication (see Hive replication) and Impala's event processor.
Point HMS at a stable endpoint (a virtual IP or DNS name the database failover updates) rather than a host. During failover, in-flight calls fail, the connection pool discards broken connections and reconnects, and clients retry through the settings above. Test that the pool validates connections, otherwise the first call after failover on every pooled connection fails.
Size connections for the fleet. Each HMS uses datanucleus.connectionPool.maxPoolSize (default 10) for two pools, one for object storage and one for the transaction handler, plus a separate compactor pool on the leader. The worker thread pool is much larger (hive.metastore.server.max.threads defaults to 1000), so under load threads queue for database connections long before the database is busy.
# Size the database for the metastore fleet, not for one server
hms_instances = 3
pools_per_instance = 2 # ObjectStore + TxnHandler each use datanucleus.connectionPool.maxPoolSize
max_pool_size = 30 # raised from the default of 10 for a busy cluster
compactor_pool = 5 # metastore.compactor.connectionPool.maxPoolSize, leader only
headroom = 1.5 # failover reconnect storms, admin sessions, replication tools
needed = (hms_instances * pools_per_instance * max_pool_size + compactor_pool) * headroom
# (3 * 2 * 30 + 5) * 1.5 = 277.5 -> set the database max_connections to at least 300
Delegation tokens and security
On a Kerberized cluster, jobs that cannot hold a Kerberos ticket, such as Spark executors or Oozie actions, authenticate to HMS with delegation tokens. The default store keeps tokens in the memory of the server that issued them, so a token issued by one HMS is unknown to the others, and a job fails the moment its connection moves. Configure hive.cluster.delegation.token.store.class to a shared store: the database-backed DBTokenStore or the ZooKeeper store. The fully qualified class name differs between releases, so take it from your version's documentation.
Every HMS server also needs its own keytab with a hive/_HOST@REALM principal, and every client needs hive.metastore.kerberos.principal set with _HOST so it can talk to whichever server it picks. With Ranger or Sentry authorization in the metastore, each server must reach the policy service too; a server that cannot load policies should fail its health check rather than serve.
Worked example: three servers across two zones
A cluster runs HiveServer2, Spark and Impala against one metastore with 40,000 tables and ACID ingestion every five minutes. The design: three HMS servers, two in zone A and one in zone B, on a managed PostgreSQL with a synchronous standby in zone B.
Clients get all three URIs with RANDOM selection, a socket timeout of 300 seconds (the slowest observed partition listing takes 90 seconds) and a socket lifetime of 30 minutes. Servers use lock election, a pool of 30 connections each, and the database token store; the database's max_connections is 300 from the calculation above. Health checks call a cheap Thrift method such as get_all_databases rather than probing only the TCP port, because a server with a dead database pool still accepts TCP connections.
Losing one HMS drops a third of capacity; clients on it retry elsewhere within seconds. Losing zone A leaves one server and the promoted database: capacity falls, so size the survivor for peak load or accept degraded latency. If the leader is lost, another server takes the lock and compaction continues.
Failure drills
Run these on a staging cluster before trusting the design:
- Kill one HMS process during a Spark job and a HiveServer2 query; both should complete, possibly after a retry pause.
- Kill the housekeeping leader; confirm another server logs that it became leader and compaction requests keep appearing.
- Fail the database over; measure the error window, confirm no committed partitions are missing, and check that event IDs did not go backwards.
- Pause an HMS with
kill -STOPinstead of killing it; this is the hung-server case, and the socket timeout decides how long clients wait. - Restart servers one at a time and check connections rebalance within the socket lifetime.
Failure modes
| Symptom | Likely cause | First response |
|---|---|---|
| Queries hang for minutes, then succeed | HMS hung, client waits for socket timeout | Lower socket timeout; health-check with a Thrift call |
| Delta files pile up, reads slow | No server running the compaction initiator | Enable lock election; check leader logs |
| Duplicate compaction requests | Initiator on several servers | Use election, not per-host flags |
| Spark job fails with invalid token after failover | In-memory token store | Switch to DBTokenStore or ZooKeeper store |
| Partitions missing after database failover | Asynchronous replication lost commits | Use synchronous or semi-synchronous replication |
| One HMS overloaded after restarts | Infinite socket lifetime pins clients | Set a finite client socket lifetime |
| Kerberos errors behind a load balancer | Principal does not match the balancer name | Add the balancer name to the keytab, or use a URI list |
HiveServer2 has its own HA story with ZooKeeper discovery; see HiveServer2 in depth. ACID tables make the transaction and lock tables part of your availability surface; see Hive ACID transactions.
What to do next
- List every client of your metastore and confirm each has all HMS URIs, not one hostname.
- Check the database's replication mode; if it is asynchronous, plan the move to synchronous or semi-synchronous replication.
- Set
metastore.housekeeping.leader.electionexplicitly and verify exactly one server runs the compaction initiator. - Configure a shared delegation token store if the cluster uses Kerberos.
- Recompute database
max_connectionsfrom the pool sizes and the number of servers. - Set client socket timeout and lifetime deliberately, then run the five failure drills and record the recovery times.