Hive ACID is usually explained through its files: a base directory, delta directories, delete deltas, and compaction that folds them together. That picture is covered in Hive ACID architecture. This article covers the other half, the part that decides which files a query may read, which writers may proceed, and which old files the cleaner may delete. That part lives in the metastore database as a transaction protocol, and nearly every hard ACID incident is a protocol incident: a writer aborted for a conflict nobody expected, a query blocked on a lock, or a directory full of deltas the cleaner refuses to remove.

Configuration names and defaults below come from the Apache Hive transactions documentation and the metastore schema. Lock-type names and some locking defaults have changed between Hive 2, 3 and 4, so confirm them against your release with SHOW LOCKS rather than trusting any article, including this one.

Advertisement

Two kinds of identity: transaction ids and write ids

A transaction id is global. Every transaction opened against the metastore, whatever it touches, gets the next number, recorded as a row in the TXNS table with a state, a user, a host, a start time and a last-heartbeat time.

A write id is per table. When a transaction first writes to a table, the metastore allocates the table's next write id and records the mapping in TXN_TO_WRITE_ID, with the per-table counter in NEXT_WRITE_ID. Delta directory names use write ids, not transaction ids: delta_512_512 means "rows written by write id 512 of this table". This indirection, introduced in Hive 3, matters for two reasons. Files remain meaningful if the table is replicated to another cluster with a different global transaction sequence, and a reader's view of one table can be expressed compactly in that table's own numbering without listing every unrelated transaction on the cluster.

Around those two tables sit the rest of the protocol state: TXN_COMPONENTS records which database, table and partition an open transaction touches; COMPLETED_TXN_COMPONENTS keeps that record after commit so the compactor knows where work happened; HIVE_LOCKS holds locks; WRITE_SET holds what committed writers modified; and MIN_HISTORY_LEVEL or MIN_HISTORY_WRITE_ID (depending on version) track the oldest snapshot that some open transaction may still need.

The life of one transaction

Transactions only exist when hive.txn.manager is set to org.apache.hadoop.hive.ql.lockmgr.DbTxnManager and hive.support.concurrency is true. With those set, a statement such as a MERGE into an orders table goes through seven steps, numbered in the diagram.

One Hive ACID transaction, as the metastore sees itHiveServer2DbTxnManageropen_txnsTXNS row, state openallocate write idTXN_TO_WRITE_ID1 open2 per tableSnapshotvalid txn + write-id listsLocksHIVE_LOCKS: read / write / exclWrite deltasdelta_512_512 in HDFS/S3345HeartbeatTXN_LAST_HEARTBEAT; miss hive.txn.timeout = abortcommit_txnwrite-set check, then state committed6overlapping writer to the same partition committed first: abortCleanermay delete obsolete deltas only below the oldest open snapshot (MIN_HISTORY_*)7
A transaction is a set of metastore rows: opened in TXNS, mapped to a per-table write id, protected by locks and heartbeats, checked for write conflicts at commit, and remembered until no snapshot needs the files it replaced.
  1. Open. HiveServer2 calls the metastore to open a transaction; a TXNS row appears in the open state.
  2. Allocate write ids. For each ACID table the statement writes, a write id is allocated and mapped.
  3. Take a snapshot. The compiler records the valid-transaction list and, per table, a valid-write-id list. This is the moment the query's view of the world is fixed.
  4. Lock. Lock requests go into HIVE_LOCKS; if an incompatible lock is held, the request waits.
  5. Write. Tasks write new delta and delete-delta directories named with the allocated write id. Nothing is visible to anyone else yet, because no snapshot considers the write id committed.
  6. Heartbeat. While running, the client periodically updates TXN_LAST_HEARTBEAT and lock heartbeats. If no heartbeat arrives for hive.txn.timeout (default 300 seconds), the housekeeping service aborts the transaction and releases its locks.
  7. Commit or abort. Commit runs the write-set conflict check, marks the transaction committed and releases locks in one metastore transaction. Abort marks it aborted; its delta directories stay on disk until the compactor and cleaner remove them.

Atomicity comes from step 7 being a single RDBMS transaction on the metastore. The files were written earlier; commit only changes which ids readers treat as committed. That is why a Hive commit is fast regardless of how much data was written, and why the metastore database is the system's single point of truth and its most important dependency. See the Hive metastore for how to run it reliably.

Advertisement

Snapshot visibility from first principles

A reader needs to answer one question for every delta directory it finds: are these rows part of my snapshot? Listing every committed write id would be unbounded, so Hive stores the answer compactly: a high watermark, the highest write id allocated when the snapshot was taken, plus a list of exceptions below it, the ids that were still open or known to be aborted. An id is visible if and only if it is at or below the watermark and not an exception.

# Snapshot visibility, conceptually (Hive's ValidWriteIdList does this per table)
class Snapshot:
    def __init__(self, high_watermark, open_ids, aborted_ids):
        self.hwm = high_watermark        # highest write id allocated when the snapshot was taken
        self.open = set(open_ids)        # allocated but not committed at that moment
        self.aborted = set(aborted_ids)  # known aborted

    def visible(self, write_id):
        if write_id > self.hwm:
            return False                 # allocated after the snapshot: invisible, even if it commits
        return write_id not in self.open and write_id not in self.aborted

# A reader starts while write id 512 is still open and 507 was aborted:
snap = Snapshot(high_watermark=512, open_ids=[512], aborted_ids=[507])
dirs = ["base_498", "delta_499_506", "delta_507_507", "delta_508_511", "delta_512_512", "delta_513_513"]
# reads base_498, delta_499_506 and delta_508_511; skips 507 (aborted), 512 (open), 513 (future)

Two consequences follow. A write id allocated after your snapshot is invisible even if it commits a millisecond later, which gives snapshot isolation: a long query sees one consistent version from start to finish. And a write id that was open when you started stays invisible to you forever, even after it commits, because your exception list does not change mid-query. The only supported isolation level is snapshot isolation; there is no serializable mode.

The cost is that the exception list grows with the number of open and aborted transactions. Many aborted transactions make every snapshot larger and every directory listing more complicated, which is why the compactor treats a count of aborted transactions on a table as a trigger (hive.compactor.abortedtxn.threshold, default 1000): compaction rewrites the data so aborted deltas can be deleted and dropped from the lists.

Locks, and why they are not enough

The lock manager backing DbTxnManager stores locks durably in the metastore rather than ZooKeeper. The metastore's lock-type enum has four values: SHARED_READ (queries), SHARED_WRITE (older documentation calls it semi-shared), EXCL_WRITE, and EXCLUSIVE (DDL such as dropping a partition). Which write lock an INSERT, UPDATE or DELETE takes depends on the Hive version and configuration, so check with SHOW LOCKS. Reads never block writes on ACID tables; DDL blocks everything on the object it changes. Locks are taken at partition granularity for partitioned tables and table granularity otherwise, never per row.

Snapshot isolation plus shared write locks leaves a classic hole: the lost update. Two UPDATE ... SET qty = qty - 1 statements compile at the same time and take the same snapshot. Both read the old value and both write a delete delta for the same row plus a new version. Locks alone cannot prevent it, because both statements legitimately hold compatible locks. HIVE-13395 (fixed in Hive 1.3.0 and 2.1.0) closed it with write-set tracking, the standard MVCC answer. At commit, a writer records what it modified in WRITE_SET; before that, it checks whether another transaction that overlapped it in time has already committed an update or delete to the same partition. If so, the later committer is aborted: first committer wins.

Operationally, the outcome depends on the write lock your version takes. Where updates take shared write locks, both pipelines run, and the later committer fails with a write-conflict error and must be retried. Where they take an exclusive write lock, the second waits for the first. Either way the data is not corrupted, and because both locks and detection work per partition, two jobs updating different rows of the same partition still collide. Plain inserts need no read-modify-write, which is one reason append-only ingestion scales far better than update-heavy ingestion.

The stuck transaction that stalls a cluster

The cleaner deletes directories made obsolete by compaction, but only when no open transaction could still read them. It uses the minimum open transaction (tracked in MIN_HISTORY_LEVEL or MIN_HISTORY_WRITE_ID, depending on version) as its limit. That limit is cluster-wide in effect: one transaction opened at 02:00 and never closed holds the watermark at 02:00 for every table, so compaction keeps producing new bases while the obsolete deltas they replaced cannot be deleted anywhere. Storage grows, directory listings get slower, and every reader's exception list carries the old id.

Heartbeat expiry normally ends this: a crashed client stops heartbeating and is aborted after hive.txn.timeout. Hive has no BEGIN or COMMIT statements (every statement auto-commits), so the danger is not a forgotten session. It is a process that keeps heartbeating while doing no useful work: a streaming ingest agent whose writer thread is stuck while its heartbeat thread keeps running, a statement waiting indefinitely on a lock behind DDL, or a very long query. It looks alive to the metastore, so nothing times it out.

Worked incident: deltas that will not go away

The symptom: sales.orders, partitioned by day and fed by an hourly MERGE, has grown from 40 to 180 directories in yesterday's partition, and queries on it have slowed threefold. SHOW COMPACTIONS shows major compactions for the partition completing successfully every few hours, with entries stuck in ready for cleaning. Compaction is working; cleaning is not.

The cleaner is waiting on an old snapshot, so find the oldest open transaction. SHOW TRANSACTIONS, or the read-only metastore query below, shows transaction 90412 opened 31 hours ago by the ingest service account from a streaming host, with a heartbeat seconds old. The agent's upstream consumer hung a day ago, so it stopped writing rows but never closed its transaction batch, and its heartbeat thread kept the transaction alive. SHOW LOCKS shows only a shared lock on the partition it streams into, so it blocks nobody's writes, only the cleaner.

-- From Beeline: who is holding transactions and locks right now?
SHOW TRANSACTIONS;                         -- id, state, user, host, start and heartbeat times
SHOW LOCKS sales.orders EXTENDED;          -- lock type and state per table/partition, with txn id
SHOW COMPACTIONS;                          -- queue: initiated, working, ready for cleaning, failed

-- Kill a zombie once you have confirmed its owner is gone
ABORT TRANSACTIONS 90412;

-- Force a compaction on one hot partition instead of waiting for the initiator
ALTER TABLE sales.orders PARTITION (ds='2026-09-29') COMPACT 'major';
-- Read-only query against the metastore RDBMS (schema names from the Hive metastore DDL)
SELECT TXN_ID, TXN_STATE, TXN_USER, TXN_HOST, TXN_STARTED, TXN_LAST_HEARTBEAT
FROM   TXNS
WHERE  TXN_STATE = 'o'                     -- open
ORDER  BY TXN_STARTED
LIMIT  20;

Stop the agent first, or it may keep its transactions alive or open new ones; then ABORT TRANSACTIONS 90412 releases whatever is left. On its next pass the cleaner deletes the obsolete deltas across every table that had been waiting, not just orders, and storage drops by several terabytes. The permanent fix is a liveness check based on rows written, not process health, plus an alert on the age of the oldest open transaction. A threshold of a few hours suits most batch clusters; tune it to your longest legitimate job.

Configuration that shapes the protocol

PropertyDocumented defaultWhat it controls
hive.txn.managerDummyTxnManagerMust be DbTxnManager for ACID
hive.support.concurrencyfalseMust be true for ACID
hive.txn.timeout300 sHeartbeat silence after which a transaction is aborted
hive.txn.heartbeat.threadpool.size5Threads used for heartbeating
hive.compactor.initiator.onfalseRun the initiator on exactly the metastore(s) you intend
hive.compactor.worker.threads0Compaction workers; must be above 0 somewhere
hive.compactor.delta.num.threshold10Delta count that triggers minor compaction
hive.compactor.delta.pct.threshold0.1Delta-to-base size ratio that triggers major compaction
hive.compactor.abortedtxn.threshold1000Aborted transactions on a table that trigger compaction

Distributions change several of these defaults, so read the effective values from HiveServer2 and the metastore with SET property; rather than assuming the Apache defaults. Raising hive.txn.timeout to hide slow heartbeats is a common mistake: it also lengthens how long a crashed client's locks and snapshot survive.

Trade-offs against table formats

Hive's protocol is pessimistic and centralised: locks and transaction state live in one RDBMS, which makes conflicts visible early and commits cheap, but makes the metastore a throughput ceiling and a single dependency. Iceberg makes the opposite trade: optimistic concurrency with an atomic swap of a metadata pointer, no lock table, and snapshots that are explicit metadata files rather than id lists. Conflicting writers still retry, but nothing like a zombie transaction pins cleanup, because snapshot expiry is an explicit maintenance operation. If your workload is dominated by concurrent updates to the same partitions, or by many engines writing one table, compare the two designs; Hive and Iceberg covers running Iceberg tables from Hive.

What to do next

  1. Confirm DbTxnManager, concurrency, and exactly one intended initiator in your effective configuration.
  2. Add an alert on the age of the oldest open transaction and on the count of compactions stuck in ready for cleaning.
  3. Give streaming ingest agents a liveness check on rows written, so a stalled writer is restarted and its transaction batch is closed.
  4. Design update and delete pipelines so that only one writer touches a given partition at a time, and add retry on write-conflict aborts.
  5. Keep the SHOW TRANSACTIONS, SHOW LOCKS ... EXTENDED, SHOW COMPACTIONS and ABORT TRANSACTIONS runbook next to your on-call notes.
  6. Watch aborted-transaction counts per table; a steady stream of aborts is a pipeline bug, not background noise.
Key takeaway: A Hive ACID commit is a metastore row change: a global transaction id maps to per-table write ids, readers fix a snapshot as a high watermark plus open and aborted exceptions, locks mediate DML and DDL at partition granularity, and write-set tracking aborts the later of two overlapping updaters. Heartbeats end crashed transactions, but an idle live one pins the cleaner for the whole cluster, so alert on the oldest open transaction. For how compaction repays delta debt, see <a href="hive_compaction.html">Hive compaction architecture</a>.