Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
129 changes: 129 additions & 0 deletions opteryx-parquet-rewritten/README.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,129 @@
# Opteryx

Opteryx is an in-process SQL query engine. Query **planning** (parse, bind,
optimize) runs in Python; query **execution** is native (Cython/C++). It
queries Parquet directly from storage with no preloading or preprocessing,
which makes it well suited to ad hoc analytics.

For more information, visit:

- [Opteryx Documentation](https://docs.opteryx.app/)
- [Opteryx GitHub Repository](https://github.com/mabel-dev/opteryx-core)

This page benchmarks Opteryx (PyPI package `opteryx-core`) on data written by
its own Parquet writer. The load step rewrites ClickBench's split Parquet files
with Opteryx's writer (rugo), using the writer's default settings; queries then
run against the rewritten files. The data stays Parquet (converted on load, the
conversion counted as load time). Its counterpart is
`Opteryx (Parquet, partitioned)`, which reads the provided files as shipped.

### Load step

`load` runs `convert.py`, which reads each `hits_N.parquet` and writes it back
with `rugo.parquet.write_parquet` (one output file per input file, row counts
checked against the source). The codec, rows per row group and row groups per
block are set at the top of `load`. Only the files the writer produces are kept,
so `data-size` measures the rewritten dataset. Nothing is precomputed: the
rewrite changes layout and encoding, not content.

### Process model

Opteryx is an in-process engine, and this entry runs it the way it is deployed:
as a long-lived service. `start` launches `server.py`, a standard-library HTTP
wrapper that imports `opteryx` and waits; `query` posts each statement to it.
This is the same shape as the pandas/polars entries.

- `BENCH_RESTARTABLE=yes`: the driver stops the server, drops the OS page cache
and starts a fresh process before every query, so try 1 is cold. Tries 2-3
run in the same process.
- `BENCH_DURABLE=yes`: the data is Parquet on disk and nothing is loaded into
process memory.
- `start` launches the service and nothing else: no warm-up query, and the
dataset is not touched. `check` reads `/health` and runs no query.
- There is no query-result cache. What carries from try 1 to tries 2-3 is
process state: imported modules and the engine's Parquet footer and schema
caches.
- The reported time is the server-side drain of `execute_to_morsels`, the same
span the previous process-per-query entry measured. Rendering the result as
TSV happens after the clock stops.

---

## Generating Benchmark Results

### High-level Steps
1. Set up the environment.
2. Install Python and the required dependencies.
3. Download the benchmark dataset.
4. Rewrite it with Opteryx's writer (this is the load step).
5. Run the benchmark script.

### Detailed Instructions

1. **Start an AWS EC2 instance**
- OS: Ubuntu 24
- Architecture: 64-bit (x86_64 or AArch64)
- Instance Type: `c6a.4xlarge`
- Root Storage: 500 GB gp2 SSD
- Advanced Details: ensure 'EBS-optimized instance' is **disabled**.

2. **SSH into the instance** (after status checks complete):
~~~bash
ssh ubuntu@{ip}
~~~

3. **Update the package list and install Git**
~~~bash
sudo apt-get update -y
sudo apt-get install git -y
~~~

4. **Clone the ClickBench repository**
~~~bash
git clone https://github.com/ClickHouse/ClickBench
cd ClickBench/opteryx-parquet-rewritten
~~~

5. **Run the benchmark script**
~~~bash
sudo ./benchmark.sh
~~~

### Python version

`opteryx-core` publishes cp314 x86_64 and AArch64 manylinux wheels and declares
no runtime dependencies, so `install` is a single binary-wheel download with no
on-box compilation and no toolchain.

`install` pins the release (`opteryx-core==0.9.164`, the latest release when
this entry was drafted) and fails if the imported
version differs, so a published result names the build that produced it.

### Query dialect

`queries.sql` adapts queries to Opteryx's dialect. The adaptations are syntactic
— they do not change what is computed, the rows returned, or the work the engine has to do:

- **Q19, Q43** — `EventTime` is stored as an integer epoch, so it is cast
explicitly (`EventTime::TIMESTAMP[s]`) before `extract(minute FROM ...)` and
before truncation.
- **Q43** — `TRUNC(<ts>, 'minute')` rather than `DATE_TRUNC('minute', <ts>)`.
- **Q29** — the `REGEXP_REPLACE` pattern and replacement use `b''` and `r''`
literals so the backslash reference survives to the regex engine.
- **Q37-Q42** — `EventDate` comparisons cast both sides to `DATE`
(`EventDate::DATE >= '2013-07-01'::DATE`).

### Hardware coverage

The results submitted with this entry are from `c6a.4xlarge`, the canonical
machine. Results for the other instance types in the ClickBench fleet come from
the ClickBench benchmark runs, not from submitted measurements.

### Known Issues

- On the memory-constrained instances the heaviest `GROUP BY` queries do not fit
in RAM and spill to swap rather than failing. They complete, but two orders of
magnitude slower — on `c6a.xlarge` (8 GB) three queries account for more than
half the total runtime. The benchmark environment provides the 16 GB swapfile
that ClickBench's `cloud-init` configures for every system; without it these
queries would be `null` instead of slow.
12 changes: 12 additions & 0 deletions opteryx-parquet-rewritten/benchmark.sh
Original file line number Diff line number Diff line change
@@ -0,0 +1,12 @@
#!/bin/bash
export BENCH_DOWNLOAD_SCRIPT="download-hits-parquet-partitioned"
# Long-lived query server (server.py). Restartable: the driver stops it, drops
# caches and starts it again before every query, so try 1 is a cold process;
# tries 2-3 hit the same process. Durable: the data is Parquet on disk, nothing
# lives in process memory, so ./load is not re-run.
export BENCH_RESTARTABLE=yes
export BENCH_DURABLE=yes
# Concurrent-QPS test stays off for now: server.py serves one request at a
# time, so N connections would measure a queue, not throughput.
export BENCH_CONCURRENT_DURATION="${BENCH_CONCURRENT_DURATION:-0}"
exec ../lib/benchmark-common.sh
12 changes: 12 additions & 0 deletions opteryx-parquet-rewritten/check
Original file line number Diff line number Diff line change
@@ -0,0 +1,12 @@
#!/bin/bash
# Succeeds only while the server is up and serving a real opteryx-core install.
# Must FAIL when the server is down: with BENCH_RESTARTABLE=yes the driver waits
# for that before dropping caches for each cold try.
#
# Runs no query — the driver calls this before every cold try, and executing
# anything here would warm the engine ahead of the measurement. The version
# shape is asserted instead of a literal so this needs no edit per release.
set -e

version=$(curl -sf http://127.0.0.1:8421/health)
[[ "$version" =~ ^[0-9]+\.[0-9]+\.[0-9]+ ]]
86 changes: 86 additions & 0 deletions opteryx-parquet-rewritten/convert.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,86 @@
#!/usr/bin/env python3
"""
Rewrite ClickBench's partitioned parquet `hits` with Opteryx's own parquet
writer (rugo), at the layout the engine reads fastest. This is the load step of
the opteryx-parquet-rewritten entry; its wall-clock is what `Load time` reports.

Self-contained on purpose: rugo and draken ship inside the `opteryx-core` wheel,
so this needs nothing but the benchmark's own install. (The engine repo's
dev/rewrite_parquet_layout.py is dev-only tooling; this is a port of its
per-file mode. Keep them in step.)

One output file per source file, same name. Each file is read whole, written
with rugo.parquet.write_parquet(compression, max_rows_per_row_group,
row_groups_per_block) and the writer's other defaults, and its row count is
checked against the source footer — a mismatch is a hard failure.

Usage: convert.py <src-dir> <dst-dir> <codec> <rows-per-row-group> <row-groups-per-block> [-j N]
codec: zstd|none (what rugo writes)

PARALLELISM: processes, not threads (morsel construction holds the GIL), three
quarters of the cores: a worker holds a whole decoded file.
"""

import os
import sys
from concurrent.futures import ProcessPoolExecutor


def rewrite(task):
src, dst, codec, rows, block = task
from draken.morsels.morsel import Morsel
from rugo.parquet import read_metadata
from rugo.parquet import read_parquet
from rugo.parquet import write_parquet

expected = read_metadata(src).num_rows
with read_parquet(src) as reader:
morsels = list(reader)
morsel = Morsel.combine(morsels) if len(morsels) > 1 else morsels[0]
data = write_parquet(
morsel, compression=codec, max_rows_per_row_group=rows, row_groups_per_block=block
)
with open(dst + ".tmp", "wb") as f:
f.write(data)
os.replace(dst + ".tmp", dst)
written = read_metadata(dst).num_rows
if written != expected or morsel.num_rows != expected:
raise RuntimeError(f"{dst}: wrote {written:,} rows, source holds {expected:,}")
return written, len(data)


def main():
argv = sys.argv[1:]
workers = max(1, (os.cpu_count() or 1) * 3 // 4)
if "-j" in argv:
i = argv.index("-j")
workers = int(argv[i + 1])
del argv[i : i + 2]
if len(argv) != 5:
print(__doc__)
return 1
src, dst, codec, rows, block = argv[0], argv[1], argv[2], int(argv[3]), int(argv[4])
if codec not in ("zstd", "none"):
print(f"ERROR: rugo writes zstd or none, not {codec!r}")
return 1
names = sorted(f for f in os.listdir(src) if f.endswith(".parquet"))
if not names:
print(f"ERROR: no parquet files in {src}")
return 1
os.makedirs(dst, exist_ok=True)
if any(f.endswith(".parquet") for f in os.listdir(dst)):
print(f"ERROR: {dst} already holds parquet files; rm -rf it first")
return 1
tasks = [(os.path.join(src, n), os.path.join(dst, n), codec, rows, block) for n in names]
with ProcessPoolExecutor(max_workers=workers) as pool:
results = list(pool.map(rewrite, tasks))
print(
f"codec={codec} rows_per_row_group={rows} row_groups_per_block={block} "
f"workers={workers} files={len(results)} rows={sum(r[0] for r in results)} "
f"bytes={sum(r[1] for r in results)}"
)
return 0


if __name__ == "__main__":
sys.exit(main())
File renamed without changes.
Original file line number Diff line number Diff line change
Expand Up @@ -30,7 +30,7 @@ fi

"$HOME/opteryx_venv/bin/python" -m pip install --upgrade pip
# Pinned so a published result is reproducible against a named release.
OPTERYX_VERSION=0.9.155
OPTERYX_VERSION=0.9.164
"$HOME/opteryx_venv/bin/python" -m pip install "opteryx-core==${OPTERYX_VERSION}"

# Fail loudly here rather than 43 queries later if pip silently fell back to a
Expand Down
23 changes: 23 additions & 0 deletions opteryx-parquet-rewritten/load
Original file line number Diff line number Diff line change
@@ -0,0 +1,23 @@
#!/bin/bash
# ClickBench ships parquet written by another writer; this entry rewrites it with
# Opteryx's own (rugo) at the layout the engine reads fastest. That rewrite IS
# the load step, and its wall-clock is what `Load time` reports.
#
# Layout: zstd (on a 500 GB gp2 volume cold runs are IO-bound, and zstd reads the
# fewest bytes; hot runs are served by the engine's decompressed chunk cache),
# 64k-row row groups, column-major blocks of 4 — rugo's writer defaults.
set -e

CODEC=zstd
ROWS_PER_ROW_GROUP=65536
ROW_GROUPS_PER_BLOCK=4

mkdir -p parquet_src
mv hits_*.parquet parquet_src/ 2>/dev/null || true

"$HOME/opteryx_venv/bin/python" convert.py parquet_src hits "$CODEC" "$ROWS_PER_ROW_GROUP" "$ROW_GROUPS_PER_BLOCK"

# Drop the source once converted: `data-size` must measure the rewritten
# dataset only.
rm -rf parquet_src
sync
File renamed without changes.
28 changes: 28 additions & 0 deletions opteryx-parquet-rewritten/query
Original file line number Diff line number Diff line change
@@ -0,0 +1,28 @@
#!/bin/bash
# Reads a SQL query from stdin and runs it on the long-lived server.
# Stdout: query result as TSV (header + rows).
# Stderr: query runtime in fractional seconds on the last line (server-measured
# drain of execute_to_morsels; see server.py).
# Exit non-zero on error; the traceback is in server.log.
set -e

body=$(mktemp)
trap 'rm -f "$body"' EXIT

meta=$(curl -sS -o "$body" -w '%{http_code} %header{x-elapsed}' \
-X POST --data-binary @- http://127.0.0.1:8421/query) || {
echo "query failed: no response from server (see server.log)" >&2
exit 1
}

code=${meta%% *}
elapsed=${meta#* }

if [ "$code" != "200" ] || [ -z "$elapsed" ]; then
echo "query failed: HTTP $code" >&2
cat "$body" >&2
exit 1
fi

cat "$body"
echo "$elapsed" >&2
59 changes: 59 additions & 0 deletions opteryx-parquet-rewritten/results/20261009/c6a.4xlarge.json
Original file line number Diff line number Diff line change
@@ -0,0 +1,59 @@
{
"system": "Opteryx (Parquet, rewritten)",
"date": "2026-10-09",
"machine": "c6a.4xlarge",
"cluster_size": 1,
"proprietary": "no",
"hardware": "cpu",
"tuned": "no",
"tags": ["C++","stateless","column-oriented","embedded"],
"load_time": 299.205,
"data_size": 9496523669,
"concurrent_qps": null,
"concurrent_error_ratio": null,
"result": [
[1.05, 0.004, 0.004],
[1.149, 0.087, 0.086],
[1.012, 0.004, 0.004],
[1.047, 0.004, 0.003],
[1.626, 0.532, 0.543],
[1.462, 0.296, 0.292],
[0.994, 0.004, 0.003],
[1.191, 0.089, 0.089],
[1.922, 0.735, 0.747],
[2.18, 0.904, 0.901],
[1.299, 0.127, 0.129],
[1.326, 0.146, 0.143],
[1.468, 0.474, 0.462],
[2.21, 0.697, 0.695],
[1.558, 0.507, 0.502],
[1.755, 0.569, 0.56],
[2.453, 1.14, 1.117],
[1.065, 0.018, 0.02],
[3.793, 2.199, 2.165],
[1.34, 0.015, 0.014],
[4.409, 0.612, 0.607],
[4.828, 0.691, 0.679],
[8.376, 1.26, 1.259],
[1.096, 0.142, 0.146],
[1.044, 0.046, 0.041],
[1.312, 0.206, 0.204],
[1.129, 0.043, 0.043],
[4.408, 0.562, 0.56],
[4.898, 2.484, 2.497],
[1.034, 0.042, 0.043],
[2.542, 0.558, 0.563],
[5.121, 0.696, 0.688],
[5.891, 3.187, 3.253],
[4.501, 2.059, 2.076],
[4.44, 2.049, 2.035],
[1.542, 0.558, 0.558],
[1.171, 0.051, 0.052],
[1.146, 0.027, 0.027],
[1.044, 0.041, 0.037],
[1.259, 0.139, 0.135],
[1.128, 0.03, 0.028],
[1.138, 0.028, 0.026],
[1.145, 0.027, 0.025]
]
}
Loading
Loading