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.

Advertisement

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.

A Trino query on Iceberg: the coordinator turns metadata into a short list of files before any worker reads dataSQL clientBI, dbt, notebookTrino coordinatorparse, plan, prunequeryCatalogREST, HMS, Gluecurrent pointermetadata.jsonschema, snapshotsManifest listper snapshotManifestsfile paths, statspartition and min/max pruningSplitsfiles, row groupsWorker 1Parquet readerWorker 2Parquet readerWorker NParquet readerObject store or HDFSdata files and delete filesWrites go the other way: workers write files, the coordinator commits a new snapshot by swapping the catalog pointer atomically
The coordinator follows the catalog pointer to the current metadata, prunes manifests and files using partition values and column statistics, and hands workers a list of splits. Commits swap the catalog pointer atomically.

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.

TypeWhere the pointer livesGood fit
hive_metastore (default)Hive metastore over ThriftExisting Hadoop estates; shared with Hive and Spark
glueAWS Glue Data CatalogAWS-native lakes; shared with Athena and EMR
restAn Iceberg REST catalog serviceMulti-engine platforms; central auth and credential vending
jdbcA relational database tableSmall or self-contained setups
nessieNessie, with branches and tagsGit-like workflows across many tables
snowflakeSnowflakeReading 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-1

For 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.

Advertisement

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

SymptomCauseFix
Queries slow down over weeksSmall files and accumulating delete filesScheduled optimize on recent partitions
Planning takes seconds before any scanToo many manifestsoptimize_manifests; fewer, larger commits
Commit failures under loadConcurrent writers on the same partitionsSeparate schedules; compact only cold partitions
Writes fail on partition countBackfill exceeds max partitions per writerBackfill in date slices
Storage keeps growingSnapshots never expired, orphan files leftexpire_snapshots and remove_orphan_files
UPDATE fails on a new tableTable created as format version 3Use version 2 until Trino supports v3 row-level writes
Users bypass access rulesDirect object-store access to table filesRestrict 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

  1. Pick the catalog type deliberately; for new multi-engine platforms, evaluate a REST catalog first.
  2. 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.
  3. Query $files and $snapshots for your three busiest tables today and record file counts, average file size and commits per day.
  4. Schedule optimize, optimize_manifests, expire_snapshots and remove_orphan_files in that order, keeping the seven-day minimum retention.
  5. Agree which engine compacts which partitions so writers and compaction never race on the same data.
  6. Alert on average data file size, delete file count and planning time, and run EXPLAIN ANALYZE on any query that regresses.
Key takeaway: Iceberg defines what a table is, the catalog decides which version is current, and Trino turns metadata into a small set of files to read or a new snapshot to commit. Most trouble comes from too many small files, delete files and manifests, so choose partitioning for size, keep row-level changes on v2, run optimize, manifest rewriting, snapshot expiry and orphan removal on a schedule, and use the metadata tables and EXPLAIN ANALYZE to see problems before your users do.