The Hive Metastore is a small database that a very large amount of data depends on. It holds the names, schemas, partitions, locations, statistics and transaction state for every table in the warehouse, and Hive, Impala, Spark, Trino and Flink all ask it where data lives. Lose it and the files are still on disk, but nothing knows what they are. Backing it up looks easy, since it is just a relational database, and that is the trap: a metastore backup is only useful if, after a restore, it agrees with the data files that kept changing after the backup was taken.
The architecture is covered in Hive Metastore and keeping the service available in Hive Metastore high availability. High availability does not replace backups: a replicated database faithfully replicates a bad DROP TABLE. This article covers taking consistent backups, restoring them, and reconciling the result with plain tables, ACID tables and Iceberg tables. Table names were checked against the Hive 4.2 schema scripts; Hive 1.x to 3.x are end-of-life, so the examples assume 4.x.
What is in the database and what is not
The metastore servers are stateless Thrift services. All state lives in the backing database, usually PostgreSQL or MySQL in production. The tables that matter most for recovery are these.
| Tables | Holds | Why it matters on restore |
|---|---|---|
| DBS, TBLS | Databases and tables, owners, table type | The catalogue itself |
| SDS, SERDES, COLUMNS_V2 | Storage descriptors: location, format, columns | Every location is an absolute URI to a filesystem |
| PARTITIONS, PARTITION_KEY_VALS | One row per partition with its own storage descriptor | Often millions of rows; partitions added after the backup disappear |
| TABLE_PARAMS, PARTITION_PARAMS | Key-value properties | Iceberg's metadata_location lives here |
| TAB_COL_STATS, PART_COL_STATS | Column statistics | Stale statistics change query plans, not results |
| TXNS, TXN_TO_WRITE_ID, NEXT_WRITE_ID | Transactions and per-table write IDs | Must never move backwards relative to the data |
| NOTIFICATION_LOG | Event stream of metadata changes | Replication and Impala's event processor track their position in it |
| VERSION | Schema version | Must match the metastore binaries |
What is not in the database is the data. Table and partition files live in HDFS or object storage, and Iceberg keeps its own metadata files beside the data. The metastore only points at them. That split is why the backup has to be designed around consistency between two systems, not just durability of one.
Three pieces of state that must agree
Suppose the backup is taken at 02:00 and the database is lost at 14:00. Restoring the 02:00 dump gives you a metastore that is twelve hours behind the data. Every change made in those twelve hours is now a mismatch: partitions added by ingestion jobs exist as directories but not as metadata; tables created exist as files nobody can find; tables dropped have metadata pointing at deleted directories. For plain external tables this is untidy but recoverable. For transactional and Iceberg tables it can corrupt data if you reopen without reconciling, which is why the rest of this article spends time on them.
The first lever is to shrink the gap. Point-in-time recovery, a base backup plus continuously archived write-ahead log or binlog, lets you restore to a minute before the failure rather than to last night. With PostgreSQL that is WAL archiving, explained in PostgreSQL WAL in depth; with MySQL it is binary logs replayed on top of a dump.
Taking a consistent backup
A metastore backup must be transactionally consistent: a partition row must never appear without its storage descriptor. Both common databases give you this without stopping the service, as long as you use the right flags.
#!/usr/bin/env bash
# Nightly logical backup of a PostgreSQL-backed metastore, plus continuous WAL archiving
# configured separately (archive_mode=on, archive_command to your backup tool).
set -euo pipefail
STAMP=$(date -u +%Y%m%dT%H%M%SZ)
OUT=/backups/hms/metastore_${STAMP}.dump
# pg_dump takes one consistent snapshot for the whole dump; -Fc allows selective restore.
pg_dump -h hms-db -U hive -d metastore -Fc -f "$OUT"
# Record what the dump describes, so a restore can be checked against it.
psql -h hms-db -U hive -d metastore -At -c '
SELECT (SELECT "SCHEMA_VERSION" FROM "VERSION"),
(SELECT count(*) FROM "TBLS"),
(SELECT count(*) FROM "PARTITIONS"),
(SELECT max("EVENT_ID") FROM "NOTIFICATION_LOG")' > "$OUT.manifest"
sha256sum "$OUT" >> "$OUT.manifest"
# MySQL / MariaDB equivalent (InnoDB): one consistent snapshot without locking tables.
# mysqldump --single-transaction --routines --triggers metastore > metastore_${STAMP}.sqlDo not copy the database's data directory while it is running, and do not back up the metastore by exporting DDL with SHOW CREATE TABLE. The DDL export loses partitions, statistics, transaction state, privileges and functions, and recreating millions of partitions from DDL takes far longer than a restore. On a managed database such as Amazon RDS or Cloud SQL, the provider's automated backups with point-in-time recovery are the base layer; still take a periodic logical dump, because it is portable to another engine version or region and can be restored into a scratch instance for testing.
The manifest written beside each dump is cheap and useful. It records the schema version, table and partition counts and the highest notification event ID at backup time, which you need later to decide which events happened after the backup.
ACID tables: write IDs must not go backwards
Hive's transactional tables store data as base and delta directories whose names carry write ID ranges, such as delta_0000041_0000041. The metastore allocates those IDs per table from NEXT_WRITE_ID, maps them to transactions in TXN_TO_WRITE_ID and decides which deltas a reader may see from the transaction state. The details are in Hive ACID transactions.
Restore an old metastore and NEXT_WRITE_ID for a table may be lower than the highest delta already on disk. The next insert is then given a write ID that an existing directory already uses. Readers can skip committed data or see data from a transaction the restored metastore considers never to have happened, and compaction can merge the wrong files. Nothing reports an error.
The safe rule is that transactional data must be restored to the same point as the metastore, either from a filesystem snapshot taken alongside the database backup or by restoring the database with point-in-time recovery to a moment after the last committed write. If neither is possible, keep the affected tables closed to writers, find the highest write ID on disk per table, and recover by exporting the readable data into a freshly created table rather than editing transaction tables by hand.
Iceberg tables: a stale pointer strands snapshots
An Iceberg table registered in the Hive catalog has almost nothing in the metastore. Its schema, partitions and snapshots are in metadata files beside the data, and the metastore keeps one table property, metadata_location, naming the current metadata JSON file. Every commit writes a new metadata file and swaps the pointer. Background is in Hive Iceberg tables.
-- Iceberg tables keep only a pointer in the metastore.
SELECT t."TBL_NAME", p."PARAM_VALUE"
FROM "TBLS" t JOIN "TABLE_PARAMS" p ON p."TBL_ID" = t."TBL_ID"
WHERE p."PARAM_KEY" = 'metadata_location';
-- After a restore, compare each pointer with the newest metadata file in the table's
-- metadata/ directory. If the directory holds newer *.metadata.json files, the restored
-- pointer is stale: snapshots committed after the backup are invisible, and their data
-- files look like orphans to remove_orphan_files.After a restore, the pointer names whatever file was current at backup time. The table opens fine and silently shows old data. Worse, snapshots committed later are no longer referenced from the current metadata, so a scheduled orphan-file cleanup will see their data files as garbage and delete them. Before re-enabling any maintenance job, compare each pointer with the newest metadata file in the table's directory, and re-point the tables whose newer metadata is valid. Never drop an Iceberg table to re-register it, because a drop can purge its data files; register the newest metadata file under a new name and swap once verified. Iceberg's Spark procedure register_table is the documented way to register an existing metadata file in a catalog. Pause orphan-file removal and snapshot expiry until every table has been checked.
Verifying backups with a restore drill
A backup you have not restored is a hope. Restore the latest dump into a scratch database on a schedule and run the metastore's own tools against it.
# Restore drill, run against a scratch database, never the live one.
pg_restore -h scratch-db -U hive -d metastore_drill --no-owner /backups/hms/metastore_LATEST.dump
# Point a throwaway hive-site.xml at the scratch database, then:
schematool -dbType postgres -info # schema version recorded vs expected by these binaries
schematool -dbType postgres -validate # sequences, tables, columns and stored locations
# Which filesystem roots do the stored locations use? Must match the cluster you restore into.
hive --service metatool -listFSRoot
# Disaster recovery into another namenode or bucket: preview, then rewrite stored locations.
hive --service metatool -updateLocation s3a://dr-warehouse hdfs://prod-nn:8020 -dryRun
hive --service metatool -updateLocation s3a://dr-warehouse hdfs://prod-nn:8020schematool -info reports the schema version stored in the database against the version the binaries expect; a mismatch means you need the matching binaries or an upgrade, covered in the Hive upgrade guides. -validate checks the schema's sequences, tables and columns and the stored locations. metatool -listFSRoot prints the filesystem roots that locations use, and -updateLocation rewrites them, with -dryRun to preview. Compare the drill's table and partition counts with the manifest; they should match exactly.
Restore runbook
- Stop every metastore server and pause ingestion, compaction, Iceberg maintenance and replication jobs, so nothing writes while you work.
- Restore the database, to a point in time if WAL or binlogs are available, into a new instance; keep the damaged one for comparison.
- Run
schematool -infoand-validate; runmetatool -listFSRootand rewrite locations if you restored into a different cluster. - Reconcile plain partitioned tables:
MSCK REPAIR TABLE t SYNC PARTITIONSadds partition directories that exist on disk and drops metadata for ones that do not; review its output on a few tables before running it everywhere. - Check every transactional table's highest on-disk write ID against
NEXT_WRITE_ID; keep mismatched tables closed to writers. - Check every Iceberg pointer against the newest metadata file.
- Start the metastore servers, then refresh readers: run
INVALIDATE METADATAin Impala, whose catalog caches metadata and tracks its position in the notification log, and restart long-running Spark and Trino sessions. - Re-baseline replication: the restored notification log ends earlier than the replica's last applied event, so incremental replication must restart from a fresh bootstrap.
Worked example: a Tuesday afternoon
A warehouse runs on a PostgreSQL metastore with 25,000 tables and 2.4 million partitions. Nightly pg_dump runs at 02:00 and WAL is archived continuously. At 14:10 a cleanup script run against the wrong database drops 300 tables. Because the drop also deletes managed tables' data, the team first checks which tables were external; their files survived.
They restore the 02:00 base backup into a new instance and replay WAL to 14:09, one minute before the script. That limits the metadata gap to one minute instead of twelve hours. schematool -validate passes and partition counts match the previous night's manifest plus that day's growth. Two Iceberg tables had commits between 14:09 and the stop of ingestion; their pointers are updated to the newest metadata files. No transactional table was written in that minute. The managed tables' data is restored separately from filesystem snapshots, which existed for exactly this reason. Impala is invalidated and replication is re-bootstrapped. Total outage: about two hours, nearly all of it waiting for the WAL replay and the checks.
Failure modes
- The backup that never ran. Cron jobs fail silently; alert on the age of the newest dump, not on job success.
- Dump without a matching schema. Restoring a 4.2 dump under 4.0 binaries fails; keep binaries or containers for the version you back up.
- Restoring over the live database. Always restore into a new instance and switch connections.
- Orphan cleanup after restore. Iceberg maintenance or Hive's directory cleanup running before reconciliation deletes valid files.
- Forgotten caches. Impala and long-running Spark sessions keep serving pre-restore metadata until refreshed.
- Replication drift. A replica that thinks it is ahead of its source applies nothing; detect it by comparing event IDs, then re-bootstrap. Replication itself is covered in Hive replication.
What to do next
- Turn on WAL archiving or binlog retention for the metastore database and set a recovery point objective in minutes.
- Schedule a consistent logical dump with a manifest of version, counts and the highest notification ID.
- Run a weekly automated restore drill into a scratch database with schematool and metatool checks.
- Decide how transactional and managed table data is snapshotted alongside the database, and document it.
- Write down the Iceberg pointer check and make pausing maintenance jobs step one of the runbook.
- Alert on the age of the newest backup and on drill failures.