Query Engines¶
This module compares engines that query external storage with managed cloud warehouses, where storage, workload isolation, governance, and operations are bundled behind a service contract.
4:15 PM. An analyst posts in the platform channel: "I need product events, customer records from Postgres, and billing from another team's MySQL, joined, by end of day." Someone replies "sure — we'll pipeline it into a warehouse, give us two weeks." The analyst says they need it today, not in two weeks.
Is a new ETL pipeline actually the right call here?
A. Yes — always land it in a warehouse first; federation is a shortcut that costs you later. B. No — a query engine can run SQL across all three stores today, at the price of paying network and each source's worst access path on every run. C. It depends on whether this becomes a recurring query or stays a one-off.
It's B for today's ask, and C is the honest answer for next quarter — which is exactly the tension this module is built around. You already have the data: product events live as Iceberg tables on object storage, customer records live in Postgres, billing lives in another team's MySQL. Nobody wants to copy 40 TB into a fourth system first — that copy is stale the moment it lands, expensive to keep, and somebody else's on-call. A query engine is the piece that runs SQL over data it does not own. Compute is a cluster. Storage is whatever the connectors can read. That split is the entire product.
Workload this module is built around¶
The SaaS analytics platform stores:
| Store | What lives there | Typical size |
|---|---|---|
| Iceberg on S3 | events — {timestamp, customer_id, user_id, service, endpoint, latency_ms, status_code, bytes} | hundreds of millions of events/day, years of history |
| PostgreSQL | customers, plans, feature_flags | tens of millions of rows, mutated all day |
| ClickHouse | serving tables for product dashboards | hot 30–90 days |
Questions that actually get asked:
- “p95 latency for
/checkoutlast Tuesday, broken down by plan tier” — Iceberg events JOIN Postgres plans. - “which customers on the enterprise plan still hit the deprecated API?” — same join, different filter.
- “reconcile yesterday’s ClickHouse dashboard with the lake” — federation, not a second ETL.
If the question is a customer-facing dashboard at 200 QPS and 200 ms, this is the wrong module. Go to ClickHouse or Pinot. If the question is “SQL over whatever we already have,” stay here.
What this module covers¶
| Topic | What you will be able to do |
|---|---|
| Trino | Plan a federated query: coordinator, workers, connectors, splits, stages, exchanges; pushdown; join distribution; why the coordinator OOMs |
| Cloud data warehouses | Explain what BigQuery, Snowflake, and Redshift actually do with a byte and a query; choose between a managed warehouse and a lakehouse/engine stack from workload evidence |
Trino is the engine we go deep on because it is the one you will actually operate against a lakehouse. Presto, Spark SQL, and BigQuery share pieces of the same mental model — stages, shuffles, stats — but they do not federate the same way.
The split: engine vs database¶
A database owns bytes on disk. Postgres, ClickHouse, Cassandra: you INSERT, they store, they query their own files.
A query engine borrows bytes. Trino never writes a MergeTree part. It asks a connector for splits, workers read them, they shuffle intermediate rows, they return a result set, and they forget.
flowchart LR
SQL["SQL"] --> C["Coordinator\nparse / plan / schedule"]
C --> W1["Worker"]
C --> W2["Worker"]
C --> W3["Worker"]
W1 --> ICE["Iceberg / S3"]
W2 --> PG["PostgreSQL"]
W3 --> CH["ClickHouse"] That picture has three consequences you cannot negotiate:
- You still pay the network. Joining Iceberg to Postgres means workers pull Iceberg columns and Postgres rows. Federation is not free compute over a magic bus.
- The source’s access path is your access path. If Postgres cannot split a 200 GB table, one Trino worker waits on one JDBC scan. The cluster size does not fix that.
- There is no storage-side index you control, except what the source already has (Iceberg manifests, Postgres B-trees, ClickHouse marks). Trino can prune; it cannot invent a primary key on someone else’s table.
How a query actually runs¶
Forget “Trino is distributed SQL” for a moment. A query is a tree of stages connected by exchanges.
SELECT c.plan, approx_percentile(e.latency_ms, 0.95)
FROM iceberg.analytics.events e
JOIN postgres.public.customers c ON e.customer_id = c.id
WHERE e.ds = DATE '2024-06-12'
AND e.endpoint = '/checkout'
GROUP BY c.plan
Rough physical shape:
| Stage | What happens | Parallelism comes from |
|---|---|---|
| Scan Iceberg | List manifests, prune to ds=2024-06-12, read customer_id, latency_ms, endpoint | splits ≈ files / row groups |
| Scan Postgres | JDBC read of customers (hopefully WHERE pushed, often not for a join build) | usually 1–few splits |
| Join + partial agg | Redistribute or broadcast, hash join, partial GROUP BY plan | workers |
| Final agg + output | Gather to coordinator (or a single output stage) | one bottleneck if the result is huge |
Two joins of the same SQL are not the same plan. If customers is 50 MB, Trino should broadcast it and skip the shuffle of events. If customers is 80 GB and stats are missing, it may partition both sides — or worse, try to broadcast and kill a worker. Stats are not a nicety; they are the difference between a 12-second query and an incident.
Deep internals, SQL, EXPLAIN, and the coordinator-memory failure mode live in Trino.
What “fast” means here¶
| Kind of question | Honest latency | Right engine |
|---|---|---|
| Ad-hoc lake SQL, 10–500 GB scanned | seconds to a minute | Trino |
| Federated Iceberg ⨝ Postgres | seconds, dominated by the slowest source + network | Trino |
| Nightly transform Iceberg → Iceberg | minutes, fine | Trino or Spark |
| Product dashboard, 200 ms, 100 QPS | milliseconds | ClickHouse / Pinot on a serving copy |
Point lookup WHERE customer_id = 42 | milliseconds | Postgres, not Trino |
Trino’s per-query tax is real: parse, analyze, split enumeration, scheduling. That tax is noise on a 40-second scan. It is the whole bill on a 20 ms lookup.
Federation is a cost model, not a feature checkbox¶
People sell federation as “query data without ETL.” True, and incomplete.
You still:
- Pay S3 LIST + GET for every Iceberg planning cycle (coordinator / Iceberg connector).
- Pay JDBC for every Postgres split, holding connections and snapshots you may not have thought about.
- Move columns you forgot to prune if the connector cannot push projections.
- Amplify a bad predicate:
WHERE lower(email) = '…'over a 2 TB Iceberg table is a full read, federated or not.
ETL is the act of paying that cost once, into a shape the serving engine likes. Federation is paying it on every query. Both are legitimate. Mixing them up is how a “simple Trino join” becomes the most expensive query in the company.
The comparison with a serving store is ClickHouse vs Trino. The lake table that makes Iceberg scans viable is Iceberg.
Scale cliffs (preview)¶
| Scale | What you notice |
|---|---|
| 10× (1 → 10 TB lake, same cluster) | Split count and S3 GET rate; small files; coordinator planning time |
| 100× | Stats go stale; broadcast joins become time bombs; Postgres side cannot keep up |
| 1000× | You stop federating the hot path. Iceberg for history, ClickHouse/Pinot for dashboards, Trino for humans and batch |
Nothing in Trino removes the need for partitioning and columnar files. The engine can only skip what the table format and file layout make skippable.
How this sits in the platform¶
flowchart TB
subgraph ingest["Ingest"]
K[Kafka]
F[Flink / Spark]
end
subgraph lake["Lake"]
I[Iceberg]
end
subgraph serving["Serving"]
CH[ClickHouse]
P[Pinot]
end
subgraph ops["Operational"]
PG[Postgres]
end
subgraph query["Ad-hoc / federation"]
T[Trino]
end
K --> F --> I
F --> CH
K --> P
I --> T
PG --> T
CH --> T Trino is the SQL front door to the lake and to systems you do not want to copy. It is not the dashboard store, not the OLTP store, and not a replacement for Spark when the job is a 12-hour shuffle with tight memory control.
How to study this module¶
- Read Trino with one query in your head: Iceberg events ⨝ Postgres customers, filtered to one day and one endpoint.
- For every operator in EXPLAIN, name whether it runs on a worker or the coordinator, and whether it crosses the network.
- Predict: broadcast or partitioned join? Then look at what happens when
customersgrows from 50 MB to 50 GB with noANALYZE. - Only then compare to ClickHouse. If you cannot say which bytes move, you cannot choose.
Check your understanding¶
An on-call dashboard joins 90 days of Iceberg events to Postgres customers on every page load (p95 8 s, 40 QPS at 09:00). A PM asks to “just point Grafana at Trino — we already have the data.”
Decide: keep federation, cache, or ETL into a serving store. Name the first bottleneck you expect.
Answer
ETL (or a scheduled incremental job) into ClickHouse or Pinot for the dashboard. 40 QPS × 8 s is not an ad-hoc workload; it is a serving SLA. Federation still pulls Iceberg files and Postgres rows on every load, so you pay S3 and JDBC at the traffic pattern of the UI.
First bottleneck is usually repeated Iceberg scans (split planning + S3) and Postgres connection / scan load at 09:00, with coordinator queued queries behind them. Caching the Trino result (or a materialized serving table) is the stopgap; the durable fix is a copy shaped for the dashboard’s GROUP BYs. Keep Trino for the analysts who need the join today against current Postgres, not for the tile that refreshes every 30 seconds.