Optimizer and Federation

Trino's Cost-Based Optimizer and Federated Queries

The cost-based optimizer (CBO) picks the join order and, per join, whether to broadcast the smaller side (REPLICATED) or hash-partition both. Manifests carry row counts and min/max values; ANALYZE orders adds distinct-value estimates as theta sketches in a Puffin .stats file (98,045 order IDs of 100,000). BookNest's customers live in PostgreSQL 1,289 : engines/shop_pg.sh starts l1-shop-pg (PostgreSQL 18, 5,000 sample customers), and a shop.properties file with connector.name=postgresql and a JDBC URL makes it catalog shop:

federated.sql: Iceberg orders joined with PostgreSQL customersSQL
-- One query, two systems: Iceberg orders on MinIO joined with PostgreSQL's customers.
SELECT c.country, count(*) AS orders, sum(o.total) AS net_revenue
FROM iceberg.booknest.orders o
JOIN shop.public.customers c ON c.customer_id = o.customer_id
WHERE o.status <> 'cancelled' AND c.signup_date >= DATE '2025-07-01'
GROUP BY c.country ORDER BY net_revenue DESC LIMIT 4;
Output
 country | orders | net_revenue
---------+--------+-------------
 US      |    131 |     4314.57
 CA      |     50 |     1906.44
 IN      |     31 |     1071.30
 GB      |     19 |      631.93
(4 rows)

EXPLAIN (engines/fed_plan.sh) shows three decisions. The PostgreSQL scan carries constraint on [signup_date]: the predicate ran inside PostgreSQL. The matching IDs came back as a dynamic filter (customer_id = #df_420) on the Iceberg 129 scan, skipping files whose customer_id range cannot match. And the orders scan was estimated at 83,333 rows after status <> 'cancelled', five sixths of the table, because six statuses are assumed equally common; the truth is 94,110. Skew is where a CBO misjudges (SET SESSION join_distribution_type overrides it). Federate small operational tables; land large ones in the lakehouse.