Reproducible benchmark harness for Trino, derived from TPC-H
Overview
In Part 1 of my TPC-H series I covered why TPC-H suits a distributed SQL engine like Trino and what its 22 queries exercise. In Part 2 I loaded the data into Iceberg tables on MinIO, ran the 22 queries over JDBC and wrote per-query statistics to CSV.
That gets you a number. It does not get you a number you can trust next month, on a different cluster, or after a Trino upgrade. This deep dive describes the harness I would build around that loop: the methodology, the environment controls, the measurement model, the automation and the statistical gate that turns a benchmark into a regression test.
It is a reference architecture, not a report of a specific engagement. Every timing in this page is illustrative.
Problem statement
A single sequential run of 22 queries answers one question: how long did these queries take, this time, on this cluster. Teams then use that answer for decisions it cannot support:
- Upgrade decisions. “Trino N+1 is 12% slower” from one run per version, when run-to-run spread on the same version is 10%.
- Sizing decisions. Worker counts chosen from a cold run that mostly measured object storage latency, or from a warm run that mostly measured a cache.
- Vendor or format comparisons. Results labelled “TPC-H” that were never audited, on hardware that was not disclosed, which breaks the TPC’s fair-use policy as well as good sense.
- Silent correctness errors. A query that returns the wrong rows quickly looks like an improvement.
The underlying problem is that the benchmark had no controlled variables, no repetition, no statistics and no definition of done.
Engineering objectives
- Reproducibility. Anyone with the repository and the environment manifest can rerun the suite and land within the documented noise band.
- Isolation of variables. One change per comparison: engine version, configuration, table layout or hardware. Everything else is pinned and recorded.
- Measurement that explains, not just times. Every run captures wall time, CPU time, peak memory, input bytes and spilled bytes, plus a plan for any query that moves.
- Statistical honesty. Report medians and percentiles with run counts. Decide regressions with a test, not by eye.
- Automation. A scripted run from data load to HTML report, and a CI job that fails a pull request on a confirmed regression.
- Fair use. Label results as a workload derived from TPC-H, disclose the environment and never report TPC primary metrics.
Solution architecture
The harness has four parts: a pinned environment, a runner, a result store and a reporting layer that feeds both humans and CI.
01 Environment
- Pinned Trino image and config
- Iceberg tables on object storage
- Catalog stats via ANALYZE
- Environment manifest
02 Runner
- Query pack pinned to spec text
- Warm-up and measured iterations
- Concurrency streams
- Result checksum validation
03 Measurement
- Client stats per query
- system.runtime.queries
- EXPLAIN ANALYZE on demand
- Node CPU, memory, IO
04 Result store
- Raw runs as Parquet
- Run metadata and git SHA
- Baselines per scale factor
05 Reporting
- Markdown and HTML report
- Statistical comparison
- CI gate with thresholds
The rule that drives the design: a result is only valid together with the manifest that produced it. The manifest records the Trino version, the image digest, config.properties, jvm.config, catalog properties, session properties, worker count and instance type, scale factor, table format and file layout, whether statistics were collected, and the git SHA of the query pack.
Technical approach
Data: generated on the fly versus materialised
Trino’s TPC-H connector generates rows deterministically at query time. That is excellent for smoke tests and plan inspection, because there is nothing to load. It is a poor basis for performance numbers: a scan of tpch.sf100.lineitem measures the generator’s CPU cost, not storage, file format or predicate pushdown.
For measured runs I materialise the data into Iceberg tables, as in Part 2, with CREATE TABLE ... AS SELECT from the connector. Two details matter more than they look:
- Column naming. By default the connector uses simplified column names without the table prefix (
orderkey, noto_orderkey). The query text in the specification uses prefixed names. Either settpch.column-naming=STANDARDin the catalog file before running the CTAS statements, or rename columns in the CTAS. Otherwise you end up quietly rewriting the queries, which changes what you are measuring. - Statistics. The cost-based optimiser chooses join order and join distribution from table statistics. Run ANALYZE on every table after loading (the Iceberg connector can also collect extended statistics on write) and record in the manifest that it was done. A run without statistics and a run with them are different experiments.
-- catalog tpch has tpch.column-naming=STANDARDCREATE SCHEMA IF NOT EXISTS lakehouse.tpch_sf100WITH (location = 's3a://bench/tpch_sf100/');
CREATE TABLE lakehouse.tpch_sf100.lineitemWITH (format = 'PARQUET', partitioning = ARRAY['month(l_shipdate)'])AS SELECT * FROM tpch.sf100.lineitem;
-- repeat for orders, customer, part, partsupp, supplier, nation, region
ANALYZE lakehouse.tpch_sf100.lineitem;Partitioning is a variable too. If you partition lineitem by ship month, say so in the report, because it changes how much data Q6 and Q14 read.
Pin the query text and validate answers
The 22 queries are templates. The specification’s qgen tool substitutes parameters and the specification publishes validation parameters with reference answers. Simplified variants of some queries circulate widely, and they produce different plans from the specification text. The harness keeps one file per query, generated once and committed, and stores a checksum of each query’s result set at the validation scale factor. A run whose checksum differs fails before its timing is even considered.
Cold, warm and repetitions
I separate three phases explicitly:
- Cold. First execution after a cluster restart, with any file system cache cleared. This is what an ad hoc analyst sees on Monday morning. Run it, report it, but never mix it with warm numbers.
- Warm-up. One or two discarded iterations per query to load metadata caches, JIT-compile hot paths and fill any file cache.
- Measured. At least five iterations per query, ideally more for short queries. I interleave the queries across iterations (Q1 to Q22, then again) rather than running Q1 five times in a row, so slow drift such as compaction or a noisy neighbour spreads across all queries instead of landing on one.
Outliers
I do not delete outliers silently. The harness keeps every observation and reports the median and p90, which are resistant to a single bad iteration. If the spread of a query (p90 divided by median) exceeds a configured limit, the report flags the query as noisy and the CI gate treats it as inconclusive rather than passing it.
Measuring per query
The runner records two views of time. The client view is wall time from submit until the last row is fetched. Draining the result set matters, as Part 2 did with its loop over next(); a query that is not fully consumed is not finished. The server view comes from the statement statistics Trino returns to the client: elapsed, queued, CPU and wall time, processed rows and bytes, physical input bytes, peak memory and spilled bytes.
CPU time is the quieter signal. On a shared cluster, wall time absorbs queueing and contention, while CPU time mostly tracks the work the plan does. I gate regressions on wall time but always show CPU time next to it, because a change in both means the plan changed, while a change in wall time alone usually means the environment did.
import time, hashlib, json, pathlib, statisticsimport trino
QUERIES = sorted(pathlib.Path("queries").glob("q*.sql"))WARMUP, ITERATIONS = 2, 7
def connect(schema, session): return trino.dbapi.connect( host="trino.bench.internal", port=8080, user="bench", catalog="lakehouse", schema=schema, session_properties=session, )
def run_once(conn, sql): cur = conn.cursor() start = time.perf_counter() cur.execute(sql) digest = hashlib.sha256() for row in cur.fetchall(): # drain fully; an unread result is not a finished query digest.update(repr(row).encode()) wall_ms = (time.perf_counter() - start) * 1000 s = cur.stats return { "query_id": s.get("queryId"), "client_wall_ms": round(wall_ms, 1), "elapsed_ms": s.get("elapsedTimeMillis"), "queued_ms": s.get("queuedTimeMillis"), "cpu_ms": s.get("cpuTimeMillis"), "peak_memory_bytes": s.get("peakMemoryBytes"), "physical_input_bytes": s.get("physicalInputBytes"), "spilled_bytes": s.get("spilledBytes"), "result_sha256": digest.hexdigest(), }
def main(run_id, schema="tpch_sf100", session=None): conn = connect(schema, session or {}) out = pathlib.Path(f"results/{run_id}.jsonl").open("w") for phase, n in (("warmup", WARMUP), ("measured", ITERATIONS)): for i in range(n): for q in QUERIES: # interleave queries across iterations rec = run_once(conn, q.read_text().rstrip(";\n")) rec.update(run_id=run_id, phase=phase, iteration=i, query=q.stem) out.write(json.dumps(rec) + "\n") out.close()The JSON lines file is the raw record. A small loader converts it to Parquet in the result store alongside the manifest, so later analysis never depends on a CSV someone edited by hand.
Plans for the queries that move
Timings say that something changed. Plans say what. For any query whose median moves beyond its threshold, the harness captures EXPLAIN ANALYZE VERBOSE on both the baseline and the candidate and stores the text next to the run.
EXPLAIN ANALYZE VERBOSESELECT nation, o_year, sum(amount) AS sum_profitFROM ( SELECT n_name AS nation, extract(year FROM o_orderdate) AS o_year, l_extendedprice * (1 - l_discount) - ps_supplycost * l_quantity AS amount FROM part, supplier, lineitem, partsupp, orders, nation WHERE s_suppkey = l_suppkey AND ps_suppkey = l_suppkey AND ps_partkey = l_partkey AND p_partkey = l_partkey AND o_orderkey = l_orderkey AND s_nationkey = n_nationkey AND p_name LIKE '%green%') AS profitGROUP BY nation, o_yearORDER BY nation, o_year DESC;Diffing two plans for the same query is usually enough to find the cause: a join that flipped from replicated to partitioned, a dynamic filter that stopped pruning lineitem, or a stage that started spilling. Reading those plans is the subject of Part 3.
Resources and concurrency
Sequential runs measure latency. The specification’s throughput test runs several query streams at once, each in a different order, and that is closer to how a shared Trino cluster behaves. The harness supports streams=N, with each stream running the 22 queries in a seeded permutation. I report per-query medians per concurrency level, total CPU seconds and peak cluster memory, and I keep node-level CPU, memory, disk and network metrics from the cluster’s own monitoring for the run window.
Memory pressure deserves a specific note. Under concurrency the large aggregations and joins (Q9, Q18, Q21 are typical candidates) can hit query.max-memory-per-node. Trino’s spill to disk can relieve this, but the documentation now describes it as legacy and points to fault-tolerant execution with a task retry policy instead. Whatever you choose, record it, and treat any non-zero spilled_bytes as a finding to explain, not noise.
Reporting
Each run produces a Markdown report (for pull requests) and an HTML report (for people), from the same data. The report opens with the environment disclosure and the fair-use disclaimer, then a per-query table.
| Query | Runs | Median (s) | p90 (s) | CPU median (s) | Peak memory | Spill | Status |
|---|---|---|---|---|---|---|---|
| Q1 | 7 | 6.8 | 7.1 | 92 | 1.1 GB | 0 | stable |
| Q6 | 7 | 2.3 | 2.6 | 31 | 0.2 GB | 0 | stable |
| Q9 | 7 | 18.4 | 21.0 | 260 | 9.8 GB | 0 | noisy |
| Q18 | 7 | 14.9 | 15.6 | 198 | 12.4 GB | 0 | stable |
| Q21 | 7 | 16.2 | 17.0 | 231 | 7.5 GB | 0 | stable |
Illustrative data, not a measured result. Format of the per-query section of the report.
Median wall time per query, derived from TPC-H at SF100
| Q1 | 6.8 s | |
|---|---|---|
| Q3 | 7.9 s | |
| Q5 | 9.6 s | |
| Q6 | 2.3 s | |
| Q9 | 18.4 s | |
| Q13 | 8.7 s | |
| Q18 | 14.9 s | |
| Q21 | 16.2 s |
The shape in Figure 2 is the useful part. In most derived-workload runs a handful of join-heavy queries dominate the total, so improvements and regressions there matter far more than a 10% swing on Q6.
Regression testing in CI
- 01Build candidate image
- 02Load SF10 from snapshot
- 03Warm-up
- 04Measured runs
- 05Compare to baseline
- 06Gate and report
The CI job runs at a small scale factor on a dedicated runner pool, compares each query with its stored baseline, and fails the build only when the regression is both larger than the per-query threshold and statistically significant. A nightly job repeats the comparison at SF100 to confirm. The statistics, the thresholds and the handling of noise are covered in Part 4.
name: perf-regressionon: [pull_request]jobs: tpch-derived: runs-on: [self-hosted, perf-dedicated] steps: - uses: actions/checkout@v4 - run: ./bench/up.sh --image "$CANDIDATE_IMAGE" --workers 4 - run: python bench/runner.py --run-id "$GITHUB_SHA" --schema tpch_sf10 - run: python bench/compare.py --baseline baselines/sf10.parquet --candidate "results/$GITHUB_SHA.jsonl" --report report.md - run: cat report.md >> "$GITHUB_STEP_SUMMARY"Validation strategy
A harness is itself software, so I validate it before trusting its verdicts.
- A/A runs. Run the same build against itself ten or more times. The comparison should almost never flag a regression. If it does, the thresholds or the environment are wrong, and the false-positive rate you measure here is the one you will live with.
- Injected regressions. Deliberately degrade the candidate, for example by setting
join_distribution_typetoPARTITIONEDfor the session or disablingenable_dynamic_filtering, and confirm the harness flags the affected queries (Q9 and Q18 typically) and not the rest. - Correctness first. Every measured iteration checks the result checksum. A failed checksum invalidates the run.
- Manifest diff. Before comparing two runs, diff their manifests. If more than the intended variable changed, the comparison is refused.
Metrics and measurement
| Metric | Source | Why it is there |
|---|---|---|
| Client wall time | Runner timer, result fully drained | What a user experiences |
| Elapsed, queued time | Statement stats, system.runtime.queries |
Separates waiting from working |
| CPU time | Statement stats | Quieter signal of plan cost |
| Peak memory | Statement stats | Headroom for concurrency |
| Physical input bytes | Statement stats | Pruning and pushdown effectiveness |
| Spilled bytes | Statement stats | Memory pressure |
| p90 to median ratio | Computed | Noise per query |
| Suite geometric mean | Computed from medians | One summary without letting the slowest query dominate |
I report a geometric mean of per-query medians as the single summary number, and I deliberately avoid anything that resembles QphH. In Part 1 I described QphH as one of TPC-H’s standardised metrics; that is true for audited results. A workload derived from TPC-H must not report it.
Challenges and trade-offs
Scale factor versus feedback time. SF100 exposes join and memory behaviour that SF1 hides, but a full suite with repetitions takes too long for every pull request. The compromise is SF10 in CI for fast detection and SF100 nightly for confirmation, accepting that some memory-bound regressions only show at night.
Object storage noise. A lakehouse on S3-compatible storage adds latency variance you do not control. Pinning a dedicated MinIO or bucket for benchmarks, and running warm, reduces it. A file system cache reduces it further but means you are no longer measuring storage. Decide which question you are asking and disclose the cache setting.
Generated versus materialised data. The connector is convenient and deterministic. Materialised tables are realistic. I use the connector for plan inspection and smoke tests and materialised Iceberg tables for every number that leaves the team.
Strictness versus noise. Tight thresholds catch small regressions and cry wolf. Loose ones are quiet and miss real problems. Per-query thresholds calibrated from A/A runs are the compromise, and noisy queries are reported as inconclusive rather than forced into pass or fail.
One environment versus many. Results from one cluster shape do not transfer linearly to another. I report results per environment and never extrapolate across them in the same table.
Outcome and lessons learned
Built this way, the harness changes the conversation from “is it slower?” to a statement of this kind (a hypothetical example): “Q18 is slower by a median of 14% with a confidence interval that excludes zero, its plan switched the orders join from replicated to partitioned after statistics went missing, and the fix is to run ANALYZE after the nightly load”. That sentence is the real deliverable of a benchmark.
The lessons that carry over from my wider practice in big data testing with Trino:
- Treat the environment manifest as part of the result. A number without it is an anecdote.
- Validate answers before timings.
- Report distributions, not single runs, and calibrate thresholds with A/A tests.
- Keep plans for anything that moves. They turn a regression report into an engineering task.
- Use honest names. A workload derived from TPC-H is useful precisely because it is not pretending to be an audited result.
Further reading
- TPC-H specification and tools
- TPC fair use quick reference
- Trino TPC-H connector
- Trino EXPLAIN ANALYZE
- Trino cost-based optimizations
- Measuring Performance Excellence: TPC-H Benchmark, Part 1
- Measuring Performance Excellence: Trino TPC-H Benchmark, Part 2
- Part 3: Reading Trino query plans to find bottlenecks in TPC-H derived queries
- Part 4: Turning a query engine benchmark into a performance regression test