Classic Impala tables are directories of files whose partitions are listed in the Hive Metastore. That design struggles at scale: listing thousands of partitions is slow, a multi-file insert is not atomic, readers can see half-written data, and changing the partition layout means rewriting the table. Apache Iceberg replaces the directory with a tree of metadata files that records exactly which data files make up each version of the table, and Impala can read and write Iceberg tables natively.
The format itself is explained in Apache Iceberg, the open table format. This article is about Impala's side: how the planner uses Iceberg metadata, the DDL and table properties Impala understands, what its writes and row-level operations actually produce, time travel and maintenance commands, and the operational rules for sharing tables with Spark and other engines. Syntax follows the Apache Impala documentation; features arrived across 4.x releases, so check your version's documentation where a version is noted.
Why Iceberg under Impala
Three problems with directory tables disappear. Atomicity: an Iceberg commit adds a new snapshot that lists the new files, and the catalog pointer to the table's current metadata file is swapped in one step, so a reader sees all of an insert or none of it. Planning cost: the file list, partition values and per-column min and max statistics live in manifest files, so planning reads a few metadata files instead of listing directories. Hidden partitioning: partitions are derived from column values by transforms such as DAY(ts), so users filter on ts and pruning just works, with no separate partition column to remember.
You also gain snapshots, which give time travel and rollback, and partition and schema evolution without rewriting data. You pay with metadata that must be maintained: snapshots and small files accumulate and need regular cleanup, which this article covers.
How Impala plans a query over the metadata tree
An Iceberg table is a chain of pointers. The catalog, by default the Hive Metastore through Iceberg's HiveCatalog, stores the location of the current metadata.json. That file holds the schema, partition specs, properties and the list of snapshots. Each snapshot points to a manifest list, which points to manifests, and each manifest lists data files or delete files with their partition values, record counts and column statistics.
Impala's catalog daemon loads and caches this metadata like any other table's, and the coordinator uses it during planning. Predicates are evaluated against partition values first, which drops whole manifests and files, then against per-file min and max statistics, which drops more. The surviving files become scan ranges for the executors, where the Parquet reader applies row-group and page statistics again. Read Impala query plans to interpret the resulting profile; the scan node reports how many files were selected out of the total.
For format version 2 tables with row-level deletes, data files that have associated position delete files are read together with those deletes, and rows whose file path and position appear in a delete file are removed by an anti-join. Files without deletes are scanned directly. The more delete files accumulate, the more of the table takes the slower path, which is the main reason compaction matters.
Iceberg statistics in manifests help pruning but say nothing about join cardinality, so run COMPUTE STATS on Iceberg tables as on any other table so the planner can order joins and size hash tables.
Creating tables: catalogs, transforms and properties
Create an Iceberg table with STORED AS ICEBERG. Partitioning uses PARTITIONED BY SPEC with transforms IDENTITY, BUCKET, TRUNCATE, YEAR, MONTH, DAY, HOUR and VOID.
CREATE TABLE analytics.events (
event_id STRING,
user_id BIGINT,
event_type STRING,
ts TIMESTAMP,
payload STRING
)
PARTITIONED BY SPEC (DAY(ts))
STORED AS ICEBERG
TBLPROPERTIES (
'format-version' = '2',
'write.parquet.compression-codec' = 'ZSTD'
);
-- CTAS works too
CREATE TABLE analytics.daily_users STORED AS ICEBERG AS
SELECT CAST(ts AS DATE) AS d, COUNT(DISTINCT user_id) AS users
FROM analytics.events GROUP BY 1;The catalog is chosen with the iceberg.catalog property: HiveCatalog is the default, and HadoopCatalog (with iceberg.catalog_location), HadoopTables and catalogs configured in hive-site.xml are also supported. Use the Hive catalog unless you have a strong reason not to; it is what Hive, Spark and Impala can all agree on, and its commit is the pointer swap the other engines expect. An EXTERNAL table with a LOCATION registers an existing Iceberg table; external.table.purge controls whether DROP TABLE deletes the data.
Set 'format-version'='2' if you will ever delete or update rows, because row-level operations require version 2. Choose partition granularity by data volume per partition: DAY(ts) for tens of gigabytes a day, MONTH for small tables, HOUR only for very high volume. BUCKET(n, col) spreads a high-cardinality key, but note that Impala does not allow INSERT OVERWRITE on tables with a bucket transform.
Writing: inserts, overwrites and the commit
Impala reads Iceberg data files in Parquet, ORC and Avro, but writes only Parquet. A table whose write.format.default is ORC can be read by Impala but not written by it, so check the property on tables created by other engines.
INSERT INTO appends: executors write Parquet files into the partition locations, and at the end of the query the new files are committed as one snapshot through the catalog. INSERT OVERWRITE replaces the partitions the query writes to. Because the commit is a swap of the metadata pointer, a failed or cancelled insert leaves only unreferenced files behind, never a half-visible table.
Iceberg uses optimistic concurrency. If another writer commits between the moment Impala read the table metadata and the moment it commits, Iceberg retries by checking whether the two changes conflict; non-conflicting appends usually succeed, while conflicting changes fail the query, which you should then rerun. Each insert writes at least one file per partition per writing fragment, so frequent small inserts produce many small files, the problem covered in small files in Hive and Impala. Batch your inserts and compact regularly.
Row-level changes: DELETE, UPDATE and MERGE
On format version 2 tables Impala supports DELETE (from Impala 4.3) and UPDATE (from 4.4), and recent releases add MERGE. All of them use merge-on-read: instead of rewriting data files, Impala writes position delete files that list the file path and row position of each removed row, and an update is a delete plus an insert of the new row version. Impala writes position deletes only; it does not produce equality deletes.
-- Erase one user's events (for example a privacy request)
DELETE FROM analytics.events WHERE user_id = 42;
-- Correct mislabelled events
UPDATE analytics.events SET event_type = 'purchase'
WHERE event_type = 'purchse' AND ts >= '2026-09-01';
-- Apply changes from a staging table
MERGE INTO analytics.user_profile
USING staging.profile_changes s ON analytics.user_profile.user_id = s.user_id
WHEN MATCHED THEN UPDATE SET email = s.email, updated_at = s.updated_at;MERGE also accepts conditional clauses that insert unmatched rows or delete matched ones; check the exact clause syntax in your release's documentation, since it has grown across versions. Merge-on-read makes writes cheap and reads more expensive, since every read must apply the deletes. A table that receives a steady stream of small deletes degrades steadily until compacted. If another engine such as Flink writes equality deletes into a table Impala reads, check your Impala version's documentation for its support of reading them before relying on it.
Time travel, history and rollback
Every commit creates a snapshot, and unexpired snapshots can be queried directly:
DESCRIBE HISTORY analytics.events;
DESCRIBE HISTORY analytics.events FROM '2026-09-30 00:00:00';
SELECT COUNT(*) FROM analytics.events FOR SYSTEM_TIME AS OF '2026-09-30 06:00:00';
SELECT COUNT(*) FROM analytics.events FOR SYSTEM_VERSION AS OF 6107342185113917312;
-- Undo a bad load by moving the current snapshot back
ALTER TABLE analytics.events EXECUTE ROLLBACK(6107342185113917312);DESCRIBE HISTORY returns the creation time, snapshot id, parent id and whether each snapshot is an ancestor of the current one. Use time travel to reproduce a report as of yesterday, compare counts before and after a load, or debug a pipeline. Rollback does not delete anything; it makes an older snapshot current again, so it is a fast, safe undo for a bad insert, as long as that snapshot has not been expired.
Evolution: partitions and schema
Partition specs can change without rewriting data. If daily partitions became too small after volume dropped, switch future writes to monthly:
ALTER TABLE analytics.events SET PARTITION SPEC (MONTH(ts));Old files keep their old layout and new files use the new spec; planning prunes each against its own spec. Use VOID(col) to stop partitioning by a column while keeping the spec history valid. Schema changes such as adding, renaming or dropping columns are metadata operations too, because Iceberg tracks columns by id rather than name. Avro-format tables are the exception; Impala does not support schema evolution for them. Compaction later rewrites old files into the latest spec and schema.
Maintenance: compaction, expiry and metadata tables
Two jobs keep an Iceberg table healthy. Compaction merges small files and folds delete files into data files. Impala provides OPTIMIZE TABLE, which rewrites the table into well-sized files, applies outstanding deletes, and rewrites old files to the latest schema and partition spec. With the FILE_SIZE_THRESHOLD_MB option it limits the rewrite to small files, which makes it cheap enough to run daily on large tables. Snapshot expiry removes old snapshots and the files only they reference:
OPTIMIZE TABLE analytics.events FILE_SIZE_THRESHOLD_MB=100;
ALTER TABLE analytics.events EXECUTE expire_snapshots('2026-09-24 00:00:00');
-- Inspect what is going on
SHOW METADATA TABLES IN analytics.events;
SELECT * FROM analytics.events.snapshots;
SELECT file_path, record_count, file_size_in_bytes FROM analytics.events.files;
SELECT COUNT(*) FROM analytics.events.delete_files;Expiry bounds your time-travel window, so choose the retention deliberately: seven days covers most debugging and rollback needs. Expiry does not remove orphan files, files that no snapshot ever referenced, such as those left by failed writes; run Iceberg's orphan file removal from Spark, for example CALL catalog.system.remove_orphan_files(table => 'analytics.events'), with an age threshold longer than your longest running write.
Watch three numbers from the metadata tables: average data file size, number of delete files, and number of snapshots. Rising delete-file counts mean reads are slowing; falling average file size means writes are too frequent or too narrow.
Sharing tables with Spark and other engines
A common layout is Spark or Flink writing and Impala serving interactive queries. The engines agree through the catalog: each commit swaps the metadata pointer in the Hive Metastore. Impala, however, caches metadata in catalogd, so after another engine commits, Impala may keep reading the old snapshot until the table is refreshed, either by event-based metadata sync if your deployment enables it, or by REFRESH analytics.events at the end of the external job. The Impala catalog article explains the caching and invalidation model.
Other rules: keep Impala-written tables on Parquet; be careful with properties one engine sets and another ignores, such as write modes for deletes; schedule compaction from one engine only, to avoid conflicting rewrites; and make sure every engine uses compatible Iceberg library versions. The Spark side of the same tables is described in Spark and Iceberg.
Failure modes
- Stale reads. Spark committed, Impala serves the old snapshot because nobody refreshed it.
- Write rejected. Impala cannot write a table whose default file format is ORC or Avro.
- Slow reads after deletes. Thousands of position delete files from frequent small deletes; fix with
OPTIMIZE TABLE. - Small-file explosion. Per-minute inserts into daily partitions; batch and compact.
- Lost rollback point. Snapshots expired too aggressively, so the bad load cannot be undone.
- Storage leak. No expiry and no orphan cleanup; the table's storage grows without bound.
- Commit conflicts. Concurrent overwrites or compactions from two engines; one fails and must be retried.
- Overwrite blocked.
INSERT OVERWRITEagainst a bucket-partitioned table.
Trade-offs
Iceberg buys atomic commits, cheap planning, hidden partitioning, time travel and evolution at the price of metadata to maintain and a dependency on catalog consistency across engines. Merge-on-read makes Impala's deletes and updates fast to write but taxes every read until compaction; copy-on-write would do the reverse, and Impala does not write it. Fine partitions prune more but create more small files. Long snapshot retention gives a deeper undo history but keeps more storage. Impala is excellent at low-latency reads over well-maintained Iceberg tables; heavy streaming ingestion and large rewrites are usually better done in Spark or Flink, with Impala serving queries.
What to do next
- Create new analytical tables as
STORED AS ICEBERGwith'format-version'='2'and a partition transform sized to daily volume. - Check that every table Impala must write has Parquet as its default write format.
- Batch inserts so each commit writes reasonably large files.
- Schedule
OPTIMIZE TABLEwith a file size threshold andexpire_snapshotswith a deliberate retention window, plus orphan file removal from Spark. - Monitor file size, delete-file count and snapshot count from the metadata tables.
- Make every external writer run
REFRESH(or rely on verified event sync) after committing. - Practise a rollback on a test table so the team knows the snapshot-id workflow before an incident.