Skip to content
Partner Developer Portal

Trino Query Optimisation

This guide covers practical Trino optimization patterns used across Explorer datasets.

Examples below are intentionally product-agnostic. Replace <catalog>, <schema>, and table names with your target product model (for example: OpenSAFELY, iPCV, IM1, Recruit, or Community Pharmacy).

This page is for anyone using Trino in Explorer to build, review, or improve SQL queries.

Use it if you are:

  • writing a new query and want good performance from the start
  • tuning an existing query that is slow or expensive to run
  • reviewing SQL and checking it follows Trino best practices
  • troubleshooting query plans using EXPLAIN or EXPLAIN ANALYZE

Use this checklist before running large or production queries:

  1. Filter early using typed predicates on key columns.
  2. Avoid CAST(column ...) in WHERE and JOIN conditions.
  3. Select only columns you need (avoid SELECT * for analytical workloads).
  4. Use partition-friendly filters (for example date or organisation keys where applicable).
  5. Prefer EXISTS for existence checks rather than LEFT JOIN ... IS NOT NULL.
  6. Validate with EXPLAIN, then confirm with EXPLAIN ANALYZE.
  7. Watch memory-heavy operations such as large window functions.
  8. Batch large result sets using partition keys or keyset pagination; never use OFFSET on billion-row tables.

Follow this workflow for reliable performance improvements:

  1. Start with a clear filter strategy.
  2. Run EXPLAIN and confirm pushdown and join strategy.
  3. Run EXPLAIN ANALYZE on a representative sample.
  4. Change one thing at a time and compare stage-level metrics.
  5. Keep the fastest version that is still easy to maintain.

Casting columns in WHERE clauses prevents the Trino planner from using column statistics and Iceberg file metadata to skip files during table scans.

SELECT *
FROM <catalog>.<schema>.observation
WHERE CAST(snomed_concept_id AS VARCHAR) = '123456789'

The planner cannot use Iceberg column statistics on snomed_concept_id (stored as BIGINT) because the predicate is CAST(snomed_concept_id AS VARCHAR) = '123456789'. Trino must read all files and perform the cast at runtime, then filter.

Impact: Full table scan; no file-level pruning.

Best practice: Compare using the native column type

Section titled “Best practice: Compare using the native column type”
SELECT *
FROM <catalog>.<schema>.observation
WHERE snomed_concept_id = BIGINT '123456789'

The predicate is on the native column type. Trino can use Iceberg column statistics to skip files that don’t contain snomed_concept_id in the target range.

Impact: File-level pruning; potentially 10–100× faster depending on table size.

WITH concept_lookups AS (
SELECT code_id, term
FROM <catalog>.<schema>.code_lookup
WHERE code_id IN (BIGINT '123456789', BIGINT '987654321')
)
SELECT obs.*
FROM <catalog>.<schema>.observation AS obs
JOIN concept_lookups AS cl
ON obs.snomed_concept_id = cl.code_id

Moving the cast to the lookup table lets the join happen on native types, preserving predicate pushdown on obs.snomed_concept_id.


When joining large tables, prefer join keys that are partition columns in Iceberg. This allows Trino to prune irrelevant partitions before the join executes.

Anti-pattern: Join on non-partition columns

Section titled “Anti-pattern: Join on non-partition columns”
SELECT pat.patient_id, obs.snomed_concept_id, obs.effective_datetime
FROM <catalog>.<schema>.patient AS pat
JOIN <catalog>.<schema>.observation AS obs
ON CAST(obs.snomed_concept_id AS VARCHAR) = pat.disease_code_lookup

Even though the join is on a meaningful key, if snomed_concept_id is not a partition column and has been cast to VARCHAR, the planner cannot prune observation partitions. A full table scan is necessary.

Best practice: Join on native partition columns when possible

Section titled “Best practice: Join on native partition columns when possible”
SELECT pat.patient_id, obs.snomed_concept_id, obs.effective_datetime
FROM <catalog>.<schema>.patient AS pat
JOIN <catalog>.<schema>.observation AS obs
ON obs.patient_id = pat.patient_id
WHERE obs.snomed_concept_id IN (BIGINT '123456789', BIGINT '987654321')
AND obs.effective_datetime >= DATE '2023-01-01'

Joining on a high-selectivity key and filtering on native typed columns helps Trino prune partitions and files where partition specs and column statistics support it.


Use EXPLAIN to inspect join strategy, data flow, and filter pushdown. Use EXPLAIN ANALYZE to see real runtime statistics.

  • TableScan details and file filters (pushdown)
  • join distribution (BROADCAST vs PARTITIONED)
  • shuffled data volume between stages
  • large estimate vs actual row mismatches
EXPLAIN
SELECT obs.patient_id, obs.snomed_concept_id
FROM <catalog>.<schema>.observation AS obs
JOIN <catalog>.<schema>.patient AS pat
ON obs.patient_id = pat.patient_id
WHERE obs.effective_datetime >= DATE '2023-01-01'

Look for:

  • Join distribution: BROADCAST (small table sent to all workers) vs PARTITIONED (data shuffled by join key)
  • Table scans: Does the plan show iceberg:fileFilter or other pushdown?
  • Unexpected cross joins: A sign of missing join conditions
EXPLAIN ANALYZE
SELECT obs.patient_id, snomed_concept_id
FROM <catalog>.<schema>.observation AS obs
WHERE obs.effective_datetime >= DATE '2023-01-01'
LIMIT 100

Provides actual row counts, memory usage, CPU time, and I/O per stage. Use this to spot:

  • Planner estimate mismatches (estimated 1M rows vs actual 100K rows)
  • High memory usage in window functions or aggregations
  • I/O or CPU bottlenecks in specific query stages

A subquery can sometimes be inlined and optimized more aggressively by the planner than a CTE. CTEs can be easier to read and maintain. Use EXPLAIN to compare behavior on your query.

CTE approach (readable, sometimes materialised)

Section titled “CTE approach (readable, sometimes materialised)”
WITH recent_observations AS (
SELECT patient_id, snomed_concept_id, effective_datetime
FROM <catalog>.<schema>.observation
WHERE effective_datetime >= DATE '2023-01-01'
)
SELECT ro.patient_id, pat.age_at_event
FROM recent_observations AS ro
JOIN <catalog>.<schema>.patient AS pat
ON ro.patient_id = pat.patient_id
LIMIT 1000;

Subquery approach (may be inlined for better pushdown)

Section titled “Subquery approach (may be inlined for better pushdown)”
SELECT ro.patient_id, pat.age_at_event
FROM (
SELECT patient_id, snomed_concept_id, effective_datetime
FROM <catalog>.<schema>.observation
WHERE effective_datetime >= DATE '2023-01-01'
) AS ro
JOIN <catalog>.<schema>.patient AS pat
ON ro.patient_id = pat.patient_id
LIMIT 1000;

Guidance: Use CTEs for readability on large queries. Use subqueries when you suspect inlining and early pushdown would reduce scanned data. Run EXPLAIN to compare plans.


Use EXISTS Over LEFT OUTER JOIN + WHERE IS NOT NULL

Section titled “Use EXISTS Over LEFT OUTER JOIN + WHERE IS NOT NULL”

For set membership or existence checks, EXISTS is more efficient than a LEFT OUTER JOIN followed by a WHERE ... IS NOT NULL filter because Trino can stop scanning the joined table once the first match is found.

Anti-pattern: LEFT OUTER JOIN to check existence

Section titled “Anti-pattern: LEFT OUTER JOIN to check existence”
SELECT DISTINCT pat.patient_id
FROM <catalog>.<schema>.patient AS pat
LEFT OUTER JOIN <catalog>.<schema>.observation AS obs
ON obs.patient_id = pat.patient_id
AND obs.snomed_concept_id = BIGINT '123456789'
WHERE obs.observation_id IS NOT NULL

Trino must scan all matching observations per patient and materialize the join before filtering on IS NOT NULL. This is expensive for high-cardinality observation tables.

SELECT DISTINCT pat.patient_id
FROM <catalog>.<schema>.patient AS pat
WHERE EXISTS (
SELECT 1
FROM <catalog>.<schema>.observation AS obs
WHERE obs.patient_id = pat.patient_id
AND obs.snomed_concept_id = BIGINT '123456789'
)

Trino implements this as a semi-join and stops scanning observations for each patient once the first match is found. Much faster for large observation tables.


Window functions like ROW_NUMBER() must materialize all input rows before computing the result. For large datasets, this consumes significant memory.

High memory usage: ROW_NUMBER on millions of rows

Section titled “High memory usage: ROW_NUMBER on millions of rows”
SELECT patient_id, observation_id, effective_datetime,
ROW_NUMBER() OVER (PARTITION BY patient_id ORDER BY effective_datetime DESC) AS rn
FROM <catalog>.<schema>.observation

On a 40M-row observation table, materializing the full window may consume 5–10 GB of memory.

Best practice: Use a filtering subquery or FETCH FIRST

Section titled “Best practice: Use a filtering subquery or FETCH FIRST”
-- Option 1: Subquery with LIMIT per patient
WITH ranked AS (
SELECT patient_id, observation_id, effective_datetime,
ROW_NUMBER() OVER (PARTITION BY patient_id ORDER BY effective_datetime DESC) AS rn
FROM <catalog>.<schema>.observation
)
SELECT *
FROM ranked
WHERE rn = 1

Run EXPLAIN to confirm Trino applies the rn = 1 predicate before materializing the full window. If the plan shows the window is computed first, then filtered, memory pressure will be high.

Alternative: Aggregate if data structure allows

Section titled “Alternative: Aggregate if data structure allows”
SELECT patient_id, MAX_BY(observation_id, effective_datetime) AS latest_observation_id
FROM <catalog>.<schema>.observation
GROUP BY patient_id

Avoids window function materialization altogether by computing the result directly in an aggregate.


Explorer datasets can contain hundreds of millions to billions of rows. Running a single unbounded query over that volume is expensive and frequently times out. Use one of the following strategies to process data incrementally and safely.

-- First batch
SELECT patient_id, observation_id, effective_datetime
FROM <catalog>.<schema>.observation
ORDER BY observation_id
LIMIT 10000 OFFSET 0;
-- Second batch
SELECT patient_id, observation_id, effective_datetime
FROM <catalog>.<schema>.observation
ORDER BY observation_id
LIMIT 10000 OFFSET 10000;

Trino must scan and sort all preceding rows to satisfy OFFSET. On a table with 500M rows, fetching page 50,000 requires materialising 500M rows before returning 10,000. Cost grows linearly with each subsequent batch.

Impact: Unbounded memory and CPU cost; query time worsens with every subsequent batch.

If the table is partitioned (for example by event_date or organisation), iterate over partition values in your application. Each batch maps directly to one or more Iceberg partitions, which Trino reads without touching other files.

-- Discover available date partitions
SELECT DISTINCT event_date
FROM <catalog>.<schema>.observation
ORDER BY event_date;
-- Process one partition at a time in your application loop
SELECT patient_id, observation_id, snomed_concept_id, effective_datetime
FROM <catalog>.<schema>.observation
WHERE event_date = DATE '2024-01-01'
ORDER BY observation_id;

Organisation (organisation) is another common partition key across Explorer datasets. Organisation codes follow the format CDB-1234.

-- Discover available organisation partitions
SELECT DISTINCT organisation
FROM <catalog>.<schema>.observation
ORDER BY organisation;
-- Process one organisation at a time
SELECT patient_id, observation_id, snomed_concept_id, effective_datetime
FROM <catalog>.<schema>.observation
WHERE organisation = 'CDB-1234'
ORDER BY observation_id;

You can also combine both partition columns to produce smaller, more targeted batches — useful when a single organisation’s data is itself very large:

SELECT patient_id, observation_id, snomed_concept_id, effective_datetime
FROM <catalog>.<schema>.observation
WHERE organisation = 'CDB-1234'
AND event_date = DATE '2024-01-01'
ORDER BY observation_id;

This is the most efficient batching strategy for analytical workloads because each batch performs a targeted file scan rather than a full table scan.

Impact: File-level pruning per batch; constant I/O cost per partition regardless of total table size.

For tables without a suitable partition column, or when you need finer-grained batch size control, track the last seen primary key and use it as the lower bound for the next query. This is sometimes called seek pagination.

-- First batch (no cursor yet)
SELECT patient_id, observation_id, effective_datetime
FROM <catalog>.<schema>.observation
ORDER BY observation_id
LIMIT 10000;
-- Subsequent batches — substitute :last_seen_id from your application
SELECT patient_id, observation_id, effective_datetime
FROM <catalog>.<schema>.observation
WHERE observation_id > :last_seen_id
ORDER BY observation_id
LIMIT 10000;

Because WHERE observation_id > :last_seen_id is a typed predicate on the native key column, Trino can use Iceberg column statistics to skip files that contain only lower values.

Tips:

  • Always include ORDER BY on the cursor column to make batches deterministic.
  • Use a partition-aligned or high-selectivity key (for example observation_id or patient_id) to maximise file pruning.
  • Store the cursor value in your application between batches, not in a Trino session variable.
  • Avoid composite cursors unless your sort order genuinely requires multiple columns.

Impact: Constant cost per batch; no rows are re-scanned or discarded.

Best practice: Modulo batching for parallel workloads

Section titled “Best practice: Modulo batching for parallel workloads”

When you need to divide a large table into a fixed number of independent slices—for example to distribute an initial bulk export across multiple workers—apply a modulo predicate on a stable integer key.

-- Split into 10 parallel shards; each worker runs a different value of MOD
-- Worker 0
SELECT patient_id, observation_id, snomed_concept_id
FROM <catalog>.<schema>.observation
WHERE MOD(observation_id, 10) = 0;
-- Worker 1
SELECT patient_id, observation_id, snomed_concept_id
FROM <catalog>.<schema>.observation
WHERE MOD(observation_id, 10) = 1;

Each shard is an independent query with no dependency on the others and can run concurrently. This is appropriate for bulk exports where row order does not matter.

Guidance: Confirm with EXPLAIN ANALYZE that the MOD predicate is applied as a filter rather than preventing file-level pruning. If observation_id is not a partition column, each shard still performs a full table scan—combine with a partition predicate to reduce I/O per shard:

SELECT patient_id, observation_id, snomed_concept_id
FROM <catalog>.<schema>.observation
WHERE event_date = DATE '2024-01-01'
AND MOD(observation_id, 10) = 0;

Impact: Horizontal scalability; independent shards run concurrently without coordination overhead.

For very large tables the most reliable approach is to combine partition-based batching with keyset pagination: the outer partition filter prunes Iceberg files to a manageable subset, and the inner keyset predicate controls batch size within that partition.

SELECT patient_id, observation_id, snomed_concept_id, effective_datetime
FROM <catalog>.<schema>.observation
WHERE event_date = DATE '2024-01-01'
AND observation_id > :last_seen_id
ORDER BY observation_id
LIMIT 50000;

PracticeBenefit
Compare columns using native types (no casting)File-level pruning via Iceberg column statistics
Join on partition columns when possiblePartition pruning during join
Use EXPLAIN ANALYZE to inspect plansSpot optimization opportunities and memory issues
Prefer EXISTS over LEFT OUTER JOIN + IS NOT NULLSemi-join optimization; lower memory usage
Minimize window function materializationLower memory pressure; faster queries
Use partition and sort columns in WHERE predicatesMaximum file-level pruning
Batch using partition keys or keyset paginationSafe, constant-cost processing of billions of rows
Avoid OFFSET on large tablesPrevents unbounded re-scan cost per batch
  • Start with small row limits while iterating (LIMIT 100 or LIMIT 1000).
  • Move to full dataset runs only after plan validation.
  • Record before-and-after runtime for important queries.
  • Keep reusable query templates for common filters and joins.
  • If performance regresses, compare EXPLAIN ANALYZE output stage by stage.

Why is my query still slow even with filters?

Section titled “Why is my query still slow even with filters?”

Common causes include type casting in predicates, filtering on non-partition columns only, or late filters applied after expensive joins.

Should I always replace CTEs with subqueries?

Section titled “Should I always replace CTEs with subqueries?”

No. Prefer readability first. Replace only when plan evidence shows a clear performance benefit.

Not always, but for large analytical queries it usually increases I/O and memory. Select only required columns for better performance.

Should I use LIMIT and OFFSET to page through large tables?

Section titled “Should I use LIMIT and OFFSET to page through large tables?”

No. OFFSET forces Trino to scan and discard all preceding rows on every batch. On billion-row tables this becomes prohibitively expensive after the first few pages. Prefer partition-based batching where a partition column exists, or keyset pagination on the primary key otherwise.