A classic Hive table is a directory. The metastore records which directories are partitions, and every query starts by listing them. Hive ACID added delta directories and a transaction manager on top of that model, which works but keeps Hive as the only engine that can safely write. An Iceberg table replaces the directory contract with a tree of metadata files, and Hive 4 ships the Iceberg integration in the box, so the same HiveQL you already run can create, write, time-travel and maintain tables that Spark, Trino, Impala and Flink also read and write.
This page is about the Hive side specifically. The format itself is covered in Apache Iceberg, the open table format and the Impala engine in Impala + Iceberg tables. Here you will see what the storage handler changes, how a Hive write commits, which DML works and where, a worked orders table from creation to rollback, the maintenance jobs you must schedule, and the failure modes that catch teams in their first month.
What changes when Hive stores a table as Iceberg
Three things move. First, the truth about which files belong to the table moves from the file system to Iceberg metadata. The metastore entry still exists, so Beeline, Ranger and catalog browsers still see the table, but its one important job is to hold the metadata_location parameter naming the current metadata JSON file.
Second, partitions stop being directories and never appear in the metastore. Hive turns PARTITIONED BY columns into Iceberg identity partitions in the table's partition spec, and a filter on updated_at prunes a table partitioned by day(updated_at) without a derived column. Iceberg calls this hidden partitioning.
Third, every INSERT, MERGE or compaction produces a new immutable snapshot, and readers always see one complete snapshot. That gives you time travel and rollback, and a new chore: snapshots and their files accumulate until you expire them.
The pieces and the versions you get
The integration is the org.apache.iceberg.mr.hive.HiveIcebergStorageHandler class plus Iceberg's formats and SerDe; STORED BY ICEBERG is shorthand for it. Since Iceberg 1.8.0 there is no separate Hive runtime connector: Hive 4.0.x bundles Iceberg 1.4.3, and Hive 4.1.x and 4.2.x bundle 1.9.1. That bundled version decides which table features Hive can read.
One limit shapes everything else: the Iceberg documentation states that DML works only on the Tez engine. Reads and DDL run anywhere, but INSERT, DELETE, UPDATE and MERGE need hive.execution.engine=tez, so jobs pinned to MapReduce or Hive on Spark cannot write Iceberg tables.
Creating tables
The basic statement is CREATE EXTERNAL TABLE ... STORED BY ICEBERG. A plain CREATE TABLE works only if the cluster runs the metastore's MetaStoreMetadataTransformer, which converts it to an external table. The default file format is Parquet. Add STORED AS ORC or STORED AS AVRO to change it. Use the Iceberg spec syntax to partition by transforms:
-- Hive 4: the Iceberg table is a Hive table whose storage handler is Iceberg.
CREATE EXTERNAL TABLE sales.orders (
order_id BIGINT,
customer_id BIGINT,
status STRING,
amount DECIMAL(12,2),
updated_at TIMESTAMP
)
PARTITIONED BY SPEC (day(updated_at), bucket(16, customer_id))
STORED BY ICEBERG
STORED AS PARQUET
TBLPROPERTIES ('format-version'='2');
DESCRIBE sales.orders; -- shows "# Partition Transform Information": updated_at DAY, customer_id BUCKET[16]The transforms match Spark's: year, month, day and hour on timestamps, plus bucket(N, col) and truncate(L, col). Choose the grain from file sizes: 40 MB a day across 16 day buckets gives 2.5 MB files, the small-file problem again. Use fewer buckets or a month grain there.
CREATE TABLE AS SELECT creates the table when the query starts and commits the data when it finishes, so for a while the table exists and is empty. Do not point a downstream job at a CTAS target until the statement returns.
Catalogs: who owns the pointer
Hive has one global catalog, the metastore it is configured with. Iceberg supports several catalog types, and a Hive table entry decides how its Iceberg table is loaded through the iceberg.catalog table property:
| iceberg.catalog | Where the table lives | Typical use |
|---|---|---|
| not set | HiveCatalog on this cluster's metastore | The default. HMS holds the pointer; Spark and Trino use the same HMS. |
| a registered name | A custom catalog defined by iceberg.catalog.NAME.* settings | Overlay tables that live in another HMS, a Hadoop catalog or AWS Glue. |
| location_based_table | A path, with no catalog | Reading a HadoopTables directory created by another tool. |
For HiveCatalog tables, Iceberg table properties and Hive table properties are kept in sync, so ALTER TABLE ... SET TBLPROPERTIES from Hive changes the Iceberg table. For other catalogs they are not, so set Iceberg properties through the owning engine.
Overlays carry one serious hazard. Create one over an existing custom-catalog table without EXTERNAL and Hive treats it as managed, so a later DROP TABLE can delete the shared data files. Always overlay with CREATE EXTERNAL TABLE.
How a Hive write commits
Follow one INSERT from start to finish. HiveServer2 compiles the statement and Tez tasks write new Parquet or ORC files under the table location. Nothing is visible yet. When the tasks finish, the commit step writes a manifest describing the new files with their column statistics, a new manifest list that includes the existing manifests plus the new one, and a new metadata JSON with the new snapshot appended. Finally it asks the metastore to change metadata_location from the version it started from to the new file, and that change only succeeds if nobody else committed in between.
This is optimistic concurrency. The second writer to swap finds that the base moved, re-checks whether its change still applies, and either retries or fails with a commit conflict. Appends on top of appends nearly always succeed on retry; two MERGE statements touching the same files conflict, so run competing row-level jobs one after another.
Atomicity is per table. A multi-insert such as FROM src INSERT INTO a ... INSERT INTO b ... commits a and then b, so a failure between them leaves only a changed. If two tables must change together, record a batch id in both and have readers check both.
Row-level DML: DELETE, UPDATE and MERGE
Hive 4 supports DELETE FROM, UPDATE and MERGE INTO on Iceberg tables, on Tez. The cost depends on what the filter matches. If a DELETE filter matches whole partitions, Iceberg drops those files from the metadata and writes no data at all. If it matches individual rows, the affected data files are rewritten (copy-on-write), which the Iceberg Hive documentation states explicitly.
Format version 2 also defines merge-on-read, where writers record deleted positions in small delete files that readers apply, selected by the Iceberg properties write.delete.mode, write.update.mode and write.merge.mode. Verify which modes your Hive build honours: run a one-row UPDATE on a test table and count rows in its delete_files table. Copy-on-write makes writes expensive and reads cheap; merge-on-read is the reverse until compaction.
MERGE has one rule that breaks pipelines: only one source row may update any given target row, or the statement fails. Change-data feeds often contain several changes to the same key, so de-duplicate to the latest change per key before the merge, as the worked example does.
Worked example: an orders table from load to rollback
The sales.orders table above receives a daily load from staging, then a feed of corrections, then one day a backfill that turns out to be wrong. Every step is plain HiveQL:
-- 1. Load a day. One INSERT is one snapshot; the commit is atomic for this table.
INSERT INTO sales.orders
SELECT order_id, customer_id, status, amount, updated_at
FROM staging.orders_raw WHERE ds = '2026-10-01';
-- 2. Apply late corrections. Only one source row may match any target row,
-- so de-duplicate the change feed first or MERGE fails.
MERGE INTO sales.orders t
USING (
SELECT * FROM (
SELECT c.*, row_number() OVER (PARTITION BY order_id ORDER BY changed_at DESC) rn
FROM staging.order_changes c) x
WHERE rn = 1) s
ON t.order_id = s.order_id
WHEN MATCHED AND s.op = 'D' THEN DELETE
WHEN MATCHED THEN UPDATE SET status = s.status, amount = s.amount, updated_at = s.changed_at
WHEN NOT MATCHED AND s.op <> 'D' THEN
INSERT VALUES (s.order_id, s.customer_id, s.status, s.amount, s.changed_at);
-- 3. Inspect what happened.
SELECT committed_at, snapshot_id, operation, summary
FROM sales.orders.snapshots ORDER BY committed_at DESC LIMIT 5;
-- 4. Tag the good state before a risky backfill, then backfill.
ALTER TABLE sales.orders CREATE TAG before_backfill;
INSERT INTO sales.orders SELECT * FROM staging.orders_backfill;
-- 5. Compare old and new without copying anything.
SELECT count(*) FROM sales.orders FOR SYSTEM_VERSION AS OF 4110952781937713041;
SELECT count(*) FROM sales.orders;
-- 6. Backfill was wrong: roll back; the tag protects that snapshot from expiry.
ALTER TABLE sales.orders EXECUTE ROLLBACK('2026-10-02 06:00:00');Step 3 uses a metadata table. Hive exposes Iceberg's metadata tables as db.table.NAME, including snapshots, history, files, data_files, delete_files, manifests, partitions and refs, and you can join and filter them like any table. Step 5 reads an old snapshot with FOR SYSTEM_VERSION AS OF; FOR SYSTEM_TIME AS OF '2026-10-02 05:00:00' does the same by timestamp. Step 6 rolls back by time, and EXECUTE ROLLBACK(snapshot_id) rolls back to an exact snapshot. Rollback only moves the current pointer to an older snapshot, so it is instant and it only works if that snapshot has not been expired.
Branches go one step further. ALTER TABLE sales.orders CREATE BRANCH audit creates an independent line of snapshots. You can write to it with INSERT INTO sales.orders.branch_audit, validate it with SELECT ... FROM sales.orders.branch_audit, and then publish it with ALTER TABLE sales.orders EXECUTE FAST-FORWARD 'main' 'audit' once the checks pass. That is a write-audit-publish pattern without a second table.
Evolution is metadata too: ALTER TABLE sales.orders SET PARTITION SPEC (hours(updated_at)) applies to new writes while old files keep their layout. One limit: PARTITION (...) clauses in INSERT, TRUNCATE and DROP PARTITION accept identity-partition columns only.
Maintenance you must schedule
Iceberg never deletes anything on its own. Every snapshot keeps its files reachable, streaming or frequent small inserts produce many small files, and merge-on-read deletes slow reads until they are compacted. Hive 4 can run all three jobs:
-- Nightly, from a scheduler (Airflow, Oozie, cron + beeline). Order matters.
ALTER TABLE sales.orders COMPACT 'major'; -- or: OPTIMIZE TABLE sales.orders REWRITE DATA;
ALTER TABLE sales.orders EXECUTE expire_snapshots('2026-09-25 00:00:00.000000000');
ALTER TABLE sales.orders EXECUTE DELETE ORPHAN-FILES OLDER THAN ('2026-09-29 00:00:00.000000000');
-- Health checks you can alert on
SELECT count(*) AS data_files, avg(file_size_in_bytes) AS avg_bytes FROM sales.orders.files;
SELECT count(*) AS delete_files FROM sales.orders.delete_files;
SELECT count(*) AS snapshots FROM sales.orders.snapshots;Compact first so the rewritten state is a new snapshot. Then expire snapshots older than your time-travel window, which removes files only the expired snapshots referenced. Then delete orphan files, meaning files no metadata references, such as output from failed writes. Give the orphan cleanup a cutoff of at least a day so it cannot delete files from a write still running, and pick one owner for expiry rather than letting Spark and Hive jobs expire with different windows.
Alert on the health queries: a rising file count with a falling average size means compaction is not keeping up, a growing delete-file count means reads are paying for merge-on-read, and a snapshot count that only grows means expiry is not running.
Failure modes
- Writes fail on a non-Tez engine. DML works only on Tez. Pin
hive.execution.engine=tezin the jobs that write Iceberg tables. - Commit conflicts under concurrent MERGE. Two row-level jobs touching the same files will not both commit. Serialise them, or partition the work so each job touches different partitions.
- Dropped data through an overlay. A managed overlay plus DROP TABLE deletes shared files. Use EXTERNAL overlays.
- Time travel and rollback that stop working. An expiry window shorter than your recovery needs removes the snapshot you want. Set the window from your incident-recovery time, not from storage cost.
- Reader too old for the table. Hive 4.0.x carries Iceberg 1.4.3. If a newer engine writes features that version cannot read, Hive queries fail. Upgrade Hive, or agree the table features all engines must support.
- In-place migration of the wrong table.
ALTER TABLE t SET TBLPROPERTIES ('storage_handler'='org.apache.iceberg.mr.hive.HiveIcebergStorageHandler')migrates external Avro, Parquet or ORC tables by writing metadata over the existing files. It does not apply to Hive ACID tables, which need a rewrite.
Trade-offs against Hive ACID tables
| Concern | Hive ACID | Hive + Iceberg |
|---|---|---|
| Who can write safely | Hive only | Any Iceberg engine through the shared catalog |
| Planning on object stores | Directory and partition listing | Metadata files with file-level statistics |
| History | None beyond compaction | Snapshots, time travel, tags, branches, rollback |
| Upkeep | Compactor inside the metastore | Your scheduled compaction, expiry and orphan cleanup |
| Execution engine for DML | Tez | Tez |
If Hive is the only engine and ACID works, it is not broken. Iceberg pays off with several writing engines, object storage or a need for rollback. The Hive ACID architecture page explains the model you would be leaving, and the metastore architecture explains the service that still holds every Iceberg pointer.
What to do next
- Confirm your Hive version and its bundled Iceberg version, and that writing jobs run on Tez.
- Create a test twin with
CREATE TABLE ... LIKE ... STORED BY ICEBERGand run the worked example against it. - Run a one-row UPDATE and inspect
delete_filesto learn which write mode your build uses. - Choose partition transforms from measured daily volume so files land in the hundreds of megabytes.
- Decide one owner for snapshot expiry and set the window from your recovery needs.
- Schedule compaction, expiry and orphan cleanup in that order, and alert on the three health queries.
- Make every overlay EXTERNAL and restrict DROP on overlays.
- Read the small-file problem before tuning bucket counts.