| name | trino-query-optimization |
| description | Trino distributed SQL query optimization — predicate/projection/aggregation pushdown, join reordering (AUTOMATIC/ELIMINATE_CROSS_JOINS), broadcast vs partitioned joins, dynamic filtering, CBO with ANALYZE, filter-early patterns, partition pruning, avoiding SELECT *, reducing shuffle, cross-catalog query cost, session property tuning, query hints, anti-patterns for slow Trino queries |
Trino Query Optimization
When to Use
- A Trino query is slower than expected and you need to diagnose and fix it
- Reviewing query patterns before deploying to production
- Choosing the right partitioning, join order, or materialization strategy
- Tuning session properties for a specific workload class
Optimization Hierarchy (apply in order)
1. Data layout — partition pruning, sorted files, file sizing
2. Pushdown — predicate / projection / aggregation / join
3. Join strategy — broadcast vs partitioned, reorder
4. Dynamic filter — wait timeout, build-side filtering
5. Memory — spill, exchange buffer, broadcast limit
6. Parallelism — task.writer-count, task.concurrency
1. Predicate Pushdown — Filter Early
Always push filters to the earliest possible stage. Trino automatically pushes WHERE predicates into connector scans.
SELECT customer_id, SUM(amount)
FROM iceberg.gold.orders
GROUP BY customer_id;
SELECT customer_id, SUM(amount)
FROM iceberg.gold.orders
WHERE order_date >= DATE '2024-01-01'
AND status = 'completed'
GROUP BY customer_id;
Identify pushdown success in EXPLAIN:
- Predicate pushed down: no
ScanFilterProject, constraint appears in TableScan
- Predicate NOT pushed down:
Filter[...] operator appears above TableScan
2. Avoid SELECT * — Projection Pushdown
SELECT * FROM iceberg.silver.orders WHERE order_date = DATE '2024-06-01';
SELECT order_id, customer_id, amount
FROM iceberg.silver.orders
WHERE order_date = DATE '2024-06-01';
3. Aggregation Pushdown
Trino can push COUNT, SUM, MIN, MAX, AVG into JDBC connectors (PostgreSQL, MySQL):
SELECT region, COUNT(*) AS orders_count
FROM postgresql.public.orders
GROUP BY region;
Aggregation pushdown is NOT applied when:
- Expression inside function:
SUM(a * b) — compute in Trino instead
ROLLUP, CUBE, GROUPING SETS present
WHERE filter present (limitation in some connectors)
4. Join Optimization
Broadcast vs Partitioned
| Strategy | When | Config |
|---|
| Broadcast | Build side < join-max-broadcast-table-size (default 100MB) | Automatic |
| Partitioned | Both tables large | Automatic |
SELECT f.order_id, d.region_name
FROM iceberg.gold.fact_orders f
JOIN iceberg.gold.dim_region d ON f.region_id = d.region_id;
SELECT f.order_id, d.category
FROM iceberg.gold.fact_orders f
JOIN iceberg.gold.dim_product d ON f.product_id = d.product_id;
Join Reordering (CBO)
The optimizer reorders joins automatically when statistics exist. Ensure stats are current:
ANALYZE iceberg.silver.orders;
ANALYZE iceberg.gold.dim_customer;
Control join reordering:
SET SESSION join_reordering_strategy = 'AUTOMATIC';
SET SESSION join_reordering_strategy = 'ELIMINATE_CROSS_JOINS';
Syntactic Join Order (when CBO disabled)
When join_reordering_strategy=NONE, Trino loads the rightmost table into memory as the build side. Write joins largest→smallest right to left:
SELECT f.order_id, d.region_name
FROM iceberg.gold.fact_orders f
JOIN iceberg.gold.dim_region d
ON f.region_id = d.region_id;
5. Dynamic Filtering
Dynamic filters propagate build-side values to the probe side at runtime, eliminating rows before they cross the network.
SELECT f.amount, d.category
FROM iceberg.gold.fact_sales f
JOIN iceberg.gold.dim_product d ON f.product_id = d.product_id
WHERE d.category = 'Electronics';
Tune dynamic filter wait:
# etc/config.properties
dynamic-filtering.small-broadcast-max-distinct-values-per-driver=1000
dynamic-filtering.small-broadcast-max-size-per-driver=512kB
dynamic-filtering.large-broadcast-max-distinct-values-per-driver=50000
dynamic-filtering.large-broadcast-max-size-per-driver=20MB
Session property:
SET SESSION dynamic_filter_wait_timeout = '2s';
6. Partition Pruning (Iceberg)
Iceberg hidden partitions are automatically pruned when filter matches partition transform:
SELECT COUNT(*) FROM iceberg.silver.orders
WHERE order_date BETWEEN DATE '2024-01-01' AND DATE '2024-01-31';
SELECT COUNT(*) FROM iceberg.silver.orders
WHERE DATE_TRUNC('month', order_date) = DATE '2024-01-01';
SELECT COUNT(*) FROM iceberg.silver.orders
WHERE order_date >= DATE '2024-01-01' AND order_date < DATE '2024-02-01';
7. Minimize Repartitioning (Shuffle)
Every PARTITION BY, GROUP BY, JOIN, and ORDER BY causes a shuffle (exchange). Minimize exchanges:
SELECT region, SUM(amount)
FROM (
SELECT region, amount
FROM iceberg.silver.orders
GROUP BY region, amount
) sub
GROUP BY region;
SELECT region, SUM(amount)
FROM iceberg.silver.orders
GROUP BY region;
Control task parallelism:
SET SESSION task_writer_count = 4;
SET SESSION task_concurrency = 8;
8. Statistics-Based Optimization (CBO)
Run ANALYZE regularly on high-churn tables so the optimizer can:
- Choose broadcast vs partitioned join
- Reorder joins by estimated row count
- Estimate cost of stage output
ANALYZE iceberg.silver.orders;
ANALYZE iceberg.silver.orders WITH (columns = ARRAY['customer_id', 'order_date', 'status']);
SELECT column_name, row_count, distinct_values_count, null_fraction
FROM iceberg.silver."orders$partitions"
LIMIT 10;
9. Useful Session Properties
SET SESSION query_max_memory = '10GB';
SET SESSION query_max_total_memory = '20GB';
SET SESSION join_distribution_type = 'AUTOMATIC';
SET SESSION join_reordering_strategy = 'AUTOMATIC';
SET SESSION join_max_broadcast_table_size = '200MB';
SET SESSION spill_enabled = true;
SET SESSION dynamic_filter_wait_timeout = '2s';
SET SESSION exchange_compression_codec = 'LZ4';
SET SESSION prefer_partial_aggregation = true;
10. Cross-Catalog Query Best Practices
Trino fetches data from each connector independently and joins results in-memory on workers. For large cross-catalog joins:
SELECT p.product_name, o.amount
FROM postgresql.public.products p
JOIN iceberg.silver.orders o ON o.product_id = p.id
WHERE o.order_date = DATE '2024-06-01';
CREATE TABLE iceberg.silver.products_snapshot AS
SELECT id, product_name FROM postgresql.public.products;
SELECT p.product_name, o.amount
FROM iceberg.silver.products_snapshot p
JOIN iceberg.silver.orders o ON o.product_id = p.id
WHERE o.order_date = DATE '2024-06-01';
Anti-Patterns
SELECT * on wide tables — Parquet/ORC reads only requested columns; SELECT * reads all and defeats projection pushdown.
- Functions on partition columns in WHERE —
YEAR(order_date) = 2024 disables Iceberg partition pruning; use range predicates instead.
- Joining large tables across catalogs (JDBC + Iceberg) — JDBC connectors fetch rows serially; always materialize JDBC data into Iceberg before large joins.
- Stale or missing statistics — CBO falls back to heuristics (ELIMINATE_CROSS_JOINS) when no stats exist; run
ANALYZE after bulk loads.
- Too many small stages from nested CTEs — each CTE creates a separate stage; flatten deeply nested CTEs when possible.
- Forcing broadcast join for large tables — if build side exceeds
join-max-broadcast-table-size, Trino OOMs on workers; use AUTOMATIC or PARTITIONED.
ORDER BY without LIMIT — full global sort across all workers is extremely expensive in distributed SQL; always add LIMIT or use window functions.
References
- Pushdown:
trino.io/docs/current/optimizer/pushdown.html
- Cost-based optimizations:
trino.io/docs/current/optimizer/cost-based-optimizations.html
- Session properties:
trino.io/docs/current/sql/set-session.html
- Related skills:
[[trino-explain-plan-review]], [[trino-iceberg-best-practices]], [[trino-file-layout-optimization]], [[trino-memory-and-spill-tuning]]