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:
-- 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;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.