Pairing Apache Iceberg with Trino is now one of the most common ways to replace a Hive-era data lake. Iceberg turns a directory of Parquet files into a table with atomic commits, schema evolution and snapshots. Trino is a distributed SQL engine that queries those tables interactively, joins them with other sources, and can also write, update and maintain them. Together they give you a warehouse-like experience on object storage you control.
The pairing works well only if you understand which component is responsible for what. Iceberg decides what a table is at any moment; Trino decides which files a query must read and how to write new ones; the catalog decides which version of the table is current. This article explains that division of labour, shows how Trino plans a query against Iceberg metadata, what it writes for row-level changes, how to run maintenance from SQL, and how to diagnose the problems that every Iceberg-on-Trino deployment meets sooner or later. Facts about the connector follow the Trino 483 documentation; check the release you run, since the connector changes often.
Who owns what
An Iceberg table is a tree of immutable files. The top is a metadata file holding the schema, partition specs, properties and a list of snapshots. Each snapshot points to a manifest list, which points to manifests, which list data files together with partition values and per-column statistics such as row counts, null counts and minimum and maximum values. A commit writes new files and then asks the catalog to swap the table's pointer from the old metadata file to the new one; if another writer got there first, the commit fails and is retried against the new state. Because nothing is ever modified in place, readers always see a consistent snapshot.
Trino adds the engine. The coordinator looks up the table in the catalog, reads the metadata tree, prunes it using the query's predicates, and turns the surviving files into splits that workers read in parallel. For writes, workers produce data files and the coordinator performs the commit. The Spark side of the same tables, including Spark's own write modes and procedures, is covered in Spark with Iceberg; many deployments write with Spark or a streaming job and query with Trino, so both views matter.
Choosing a catalog
The Iceberg connector supports six catalog types through iceberg.catalog.type. The choice is about who else must see the same tables and how commits are made atomic.
| Type | Where the pointer lives | Good fit |
|---|---|---|
hive_metastore (default) | Hive metastore over Thrift | Existing Hadoop estates; shared with Hive and Spark |
glue | AWS Glue Data Catalog | AWS-native lakes; shared with Athena and EMR |
rest | An Iceberg REST catalog service | Multi-engine platforms; central auth and credential vending |
jdbc | A relational database table | Small or self-contained setups |
nessie | Nessie, with branches and tags | Git-like workflows across many tables |
snowflake | Snowflake | Reading Snowflake-managed Iceberg tables; writes happen in Snowflake |
# etc/catalog/lake.properties -- REST catalog on S3
connector.name=iceberg
iceberg.catalog.type=rest
iceberg.rest-catalog.uri=https://catalog.internal:8181
iceberg.rest-catalog.warehouse=analytics
iceberg.rest-catalog.security=OAUTH2
# plus the OAuth2 credential properties your catalog provider requires
fs.native-s3.enabled=true
s3.region=eu-west-1For an existing Hadoop cluster, hive_metastore with hive.metastore.uri pointing at the metastore is the shortest path, and it lets Hive see the same tables, as described in Hive and Iceberg. For a new platform with several engines, a REST catalog is usually the better long-term choice, because authorisation and storage credentials can be handled centrally instead of being configured separately in every engine.
Designing tables from Trino
Iceberg partitioning is hidden: you partition by a transform of a column, and queries filter on the column itself. Nobody needs to know that a day(event_ts) partition exists to benefit from it. Trino exposes this through table properties.
CREATE TABLE lake.web.events (
event_id varchar,
customer_id bigint,
event_type varchar,
event_ts timestamp(6) with time zone,
payload varchar
)
WITH (
format = 'PARQUET',
format_version = 2,
partitioning = ARRAY['day(event_ts)', 'bucket(customer_id, 16)'],
sorted_by = ARRAY['customer_id']
);
-- Partition evolution: new data is written with the new spec, old files keep theirs
ALTER TABLE lake.web.events SET PROPERTIES partitioning = ARRAY['hour(event_ts)', 'bucket(customer_id, 16)'];Defaults in Trino 483 are Parquet, ZSTD compression and format version 2. Version 3 is accepted but documented as experimental, and Trino does not yet support row-level updates and deletes on v3 tables, so stay on v2 for tables you will modify. Choose partition granularity so each partition holds at least a few hundred megabytes per day; a fine partitioning on a modest table produces thousands of small files and slows every query. sorted_by makes writers sort data within each file, which tightens the min and max statistics and lets the engine skip more files and row groups for selective filters.
How Trino plans an Iceberg query
Pruning happens in layers, and each layer is cheaper than the next. First, partition values recorded in the manifests let the coordinator skip whole groups of files without opening them. Second, per-file column statistics in the manifests skip individual files whose value range cannot match the predicate. Third, workers read Parquet footers and skip row groups using their own statistics. Dynamic filtering adds a fourth layer for joins: when a small dimension table is filtered, Trino collects the matching join keys and uses them to prune the large table's splits, waiting up to iceberg.dynamic-filtering.wait-timeout (one second by default) during split generation.
The coordinator caches metadata files in memory by default, which helps repeated queries. The work it cannot avoid is reading the manifests that survive partition pruning, so a table with an enormous number of small manifests makes planning slow even when the query is selective. That is why manifest rewriting is part of maintenance. To see the effect of pruning on a real query, use EXPLAIN ANALYZE and compare input rows and the number of splits with what you expected from the predicate.
EXPLAIN ANALYZE
SELECT event_type, count(*)
FROM lake.web.events
WHERE event_ts >= TIMESTAMP '2026-09-30 00:00:00 UTC'
AND event_ts < TIMESTAMP '2026-10-01 00:00:00 UTC'
AND customer_id = 4711
GROUP BY event_type;
What row-level writes produce
Trino supports INSERT, UPDATE, DELETE and MERGE on Iceberg tables. On a v2 table, a row-level change does not rewrite the data files it touches. Instead Trino writes position delete files, each listing a data file path and the row positions that are no longer live, plus new data files for updated or inserted rows. This merge-on-read approach makes writes cheap and pushes the cost to readers, who must apply the deletes while scanning. A DELETE whose predicate covers whole partitions on identity-partitioned columns can instead drop those files outright, which is much cheaper.
The consequence is predictable: a table receiving frequent small MERGE statements accumulates delete files, and read latency climbs until compaction folds the deletes into rewritten data files. Writers also need limits. Each writer task handles at most iceberg.max-partitions-per-writer partitions (100 by default) and a query that exceeds it fails, which usually signals a backfill across too many partitions at once; split the backfill by date rather than raising the limit blindly. File size is capped near iceberg.target-max-file-size, one gigabyte by default.
Concurrency is optimistic. Two writers that change overlapping data race to commit, and the loser fails and must retry. A long-running MERGE in Trino and a compaction job in Spark on the same partitions is a classic collision; schedule them apart or compact only partitions that are no longer receiving changes.
Maintenance from SQL
Every write adds files and snapshots, and nothing is removed automatically. Trino runs the essential maintenance through ALTER TABLE ... EXECUTE.
-- 1. Compact small files and fold in delete files, recent partitions only
ALTER TABLE lake.web.events EXECUTE optimize(file_size_threshold => '128MB')
WHERE event_ts >= TIMESTAMP '2026-09-24 00:00:00 UTC';
-- 2. Rewrite manifests so planning reads fewer, larger manifest files
ALTER TABLE lake.web.events EXECUTE optimize_manifests;
-- 3. Drop snapshots older than the retention window (and the files only they reference)
ALTER TABLE lake.web.events EXECUTE expire_snapshots(retention_threshold => '7d');
-- 4. Delete files that no snapshot references, e.g. from failed writes
ALTER TABLE lake.web.events EXECUTE remove_orphan_files(retention_threshold => '7d');optimize rewrites files smaller than file_size_threshold (100 MB if you omit it) and accepts a WHERE clause, so you can compact only recent partitions; keep the filter aligned with partition boundaries. expire_snapshots and remove_orphan_files refuse retention thresholds shorter than iceberg.expire-snapshots.min-retention and iceberg.remove-orphan-files.min-retention, both seven days by default. Keep that guard. Orphan removal with a short threshold can delete files that a writer has just written but not yet committed, which corrupts the next commit. Run the four steps in that order, on a schedule, and track how long each takes.
Small files are the root of most Iceberg performance problems, exactly as on HDFS; the mechanics differ but the lesson in the HDFS small-files problem carries over directly: per-file overhead dominates when files are tiny.
Diagnosing with metadata tables
Trino exposes Iceberg metadata as hidden tables you query with ordinary SQL, by quoting the table name with a suffix: $snapshots, $history, $files, $manifests, $partitions, $refs, $properties and others. They answer most operational questions without any extra tooling.
-- How many data and delete files, and how small are they?
SELECT content, count(*) AS files,
round(avg(file_size_in_bytes) / 1048576.0, 1) AS avg_mb,
sum(record_count) AS records
FROM lake.web."events$files"
GROUP BY content; -- 0 = data, 1 = position deletes, 2 = equality deletes
-- How fast are snapshots accumulating?
SELECT date_trunc('hour', committed_at) AS hour, count(*) AS commits
FROM lake.web."events$snapshots"
GROUP BY 1 ORDER BY 1 DESC LIMIT 24;
Time travel and rollback
Because old snapshots remain until expired, you can query the table as it was. FOR VERSION AS OF takes a snapshot ID or the name of a branch or tag, and FOR TIMESTAMP AS OF takes a point in time. If a bad job corrupts a table, ALTER TABLE ... EXECUTE rollback_to_snapshot(id) points the table back at an earlier snapshot without copying data. All of this works only within the snapshot retention window, which is the real trade-off when you choose retention_threshold: longer retention means longer recovery options and more storage.
SELECT count(*) FROM lake.web.events FOR TIMESTAMP AS OF TIMESTAMP '2026-09-30 06:00:00 UTC';
SELECT snapshot_id, committed_at, operation FROM lake.web."events$snapshots" ORDER BY committed_at DESC LIMIT 5;
ALTER TABLE lake.web.events EXECUTE rollback_to_snapshot(8954597067493422955);
Worked example: the dashboard that slowed down
A product team's dashboard over lake.web.events went from two seconds to forty over a month. The table is fed by a streaming job that commits every minute and by an hourly Trino MERGE that applies late corrections. The diagnosis took three queries. The $snapshots query showed about 1,500 commits a day. The $files query showed hundreds of thousands of data files averaging a few megabytes, plus tens of thousands of position delete files. EXPLAIN ANALYZE showed most of the time in split generation and in applying deletes, not in scanning.
The fixes followed the cause. The streaming job's commit interval was raised to five minutes, cutting new files by a factor of five. A nightly job now runs optimize on the previous seven days, then optimize_manifests, then expiry and orphan removal with seven-day retention. The hourly MERGE was restricted to partitions from the last two days so compaction of older partitions never conflicts with it. After a week the dashboard was back under three seconds, and the team added an alert on average data file size per table.
Failure modes
| Symptom | Cause | Fix |
|---|---|---|
| Queries slow down over weeks | Small files and accumulating delete files | Scheduled optimize on recent partitions |
| Planning takes seconds before any scan | Too many manifests | optimize_manifests; fewer, larger commits |
| Commit failures under load | Concurrent writers on the same partitions | Separate schedules; compact only cold partitions |
| Writes fail on partition count | Backfill exceeds max partitions per writer | Backfill in date slices |
| Storage keeps growing | Snapshots never expired, orphan files left | expire_snapshots and remove_orphan_files |
| UPDATE fails on a new table | Table created as format version 3 | Use version 2 until Trino supports v3 row-level writes |
| Users bypass access rules | Direct object-store access to table files | Restrict storage access; enforce policy in Trino |
The last row deserves emphasis. Trino's access control, whether file-based or through Apache Ranger, governs SQL access. The files themselves sit in a bucket or directory, and anyone with storage permissions can read them directly. Lock down storage to the engines' service identities, as you would for any data in S3.
What to do next
- Pick the catalog type deliberately; for new multi-engine platforms, evaluate a REST catalog first.
- Create tables in format version 2 with day or hour partitioning sized for hundreds of megabytes per partition, and set sorted_by on the most selective filter column.
- Query $files and $snapshots for your three busiest tables today and record file counts, average file size and commits per day.
- Schedule optimize, optimize_manifests, expire_snapshots and remove_orphan_files in that order, keeping the seven-day minimum retention.
- Agree which engine compacts which partitions so writers and compaction never race on the same data.
- Alert on average data file size, delete file count and planning time, and run EXPLAIN ANALYZE on any query that regresses.