معیارسنجی پایگاههای داده برای بارهای کاری سری زمانی مالی
خلاصه
این سند معیارسنجی موتورهای پایگاه داده را برای دادههای سری زمانی مالی توضیح میدهد. نوشتن گروهی، پویش کامل، پرسوجوهای بازه زمانی، تجمیع OHLCV و اتصالهای مبتنی بر آخرین داده موجود را برای دادههای معامله و مظنه مقایسه میکند. روششناسی، تبدیل سمت کارخواه را بخشی از هزینه نوشتن میداند، تا زمان ماندگاری و آماده پرسوجو شدن داده صبر میکند و تعداد دقیق نتایج را اعتبارسنجی میکند تا پرسوجوهای ناقص بهطور مصنوعی سریع به نظر نرسند. خواندن با گرمسازی آغاز میشود و سپس اندازهگیریهای گرم تکراری انجام میگیرد تا حافظههای نهان موجود در موتورهای تعبیهشده و سرورها بازتاب یابند.
معیارسنجی توضیح میدهد چرا چیدمان ذخیرهسازی بر عملکرد پویش اثر دارد و چرا عملکرد پرسوجوی بازهای برای بکتست اهمیت جداگانهای دارد. همچنین هشدار میدهد اندازههای دیسک گزارششده مختص هر موتورند و مستقیماً بهعنوان سنجه فشردهسازی قابل مقایسه نیستند. نتایج به موتورهای در دسترس، سختافزار، رابطهای کارخواه و وضعیت حافظه نهان بستگی دارند؛ دفترچه پایگاه دادهای را سریعترین گزینه در همه شرایط نمینامد و توصیه میکند مقایسه با مسیر ورود داده موردنظر دوباره اجرا شود. نتایج تخصصی kdb+/PyKX فقط وقتی ارائه میشوند که محیط لازم پیکربندی شده باشد.
ایدههای کلیدی
- نوشتن، پویش، پرسوجوهای بازه زمانی، تجمیعها و اتصالهای مبتنی بر آخرین داده موجود را بهعنوان بارهای کاری جداگانه معیارسنجی کنید.
- ساخت اشیای سمت کارخواه و زمان آماده پرسوجو شدن داده را در اندازهگیری نوشتن بگنجانید.
- تعداد دقیق ردیفهای نتایج پرسوجو را اعتبارسنجی کنید تا نتایج ناقص صرفاً بهدلیل سرعت برنده نشوند.
- زمانسنجی خواندن گرم، عملکرد با داده ذخیرهشده در حافظه نهان را نشان میدهد و نباید با عملکرد خواندن سرد اشتباه گرفته شود.
- موتورهای مختلف اندازه دیسک را با قواعد حسابداری متفاوت گزارش میکنند؛ بنابراین، این اندازهها مستقیماً قابل مقایسه نیستند.
برچسبها
متن کامل
# 21_storage_benchmark_database.py
```py
# ---
# jupyter:
# jupytext:
# cell_metadata_filter: tags,-all
# text_representation:
# extension: .py
# format_name: percent
# format_version: '1.3'
# jupytext_version: 1.19.3
# kernelspec:
# display_name: Python 3 (ipykernel)
# language: python
# name: python3
# ---
# %% [markdown]
# # Storage Benchmark: Database Engines
#
# **Docker image**: `benchmark`
#
# > **Docker required**: This notebook depends on the `benchmark` environment and
# > database services. Run with:
# > ```bash
# > docker compose --profile benchmark up -d timescaledb clickhouse questdb influxdb
# > docker compose --profile benchmark run --rm benchmark python 02_financial_data_universe/21_storage_benchmark_database.py
# > ```
#
# **Focus**: Query-capable databases for financial time-series
#
# ## Database Categories
#
# | Category | Examples | Characteristics |
# |----------|----------|-----------------|
# | **Embedded** | SQLite, DuckDB, ArcticDB | In-process, no server required |
# | **Time-Series Servers** | ClickHouse, QuestDB, TimescaleDB | Production-scale, Docker required |
# | **General Purpose** | PostgreSQL, InfluxDB | Baseline comparisons |
# | **HFT Specialized** | kdb+/PyKX | Industry standard for trading |
#
# ## Operations Benchmarked
#
# 1. **Write** - Bulk insert performance
# 2. **Read** - Full table scan
# 3. **Range Query** - 7-day time window (backtesting workflow)
# 4. **OHLCV Aggregation** - Resample minute bars to daily bars
# 5. **ASOF Join** - Trade-quote alignment (critical for microstructure)
#
# ## Timing Policy
#
# One policy, applied identically to every engine. Engines that were timed
# under different rules would not belong on the same chart, so the rules are
# stated here and enforced by the `time_write` / `time_read` helpers rather
# than by per-call arguments:
#
# - **Writes** (`time_write`) - a single shot, no warm-up, against a table
# created fresh for the run. Bulk load happens once per dataset, so that is
# what we time. Repeating it would append duplicate rows on the append-only
# engines, or require a teardown that is not part of the write.
# - **A write bar is what the client pays, and the clients differ.** SQLite,
# DuckDB, ClickHouse and kdb+ take the panel in one block (`to_sql`, a Parquet
# scan, `insert_df`, `set`). PostgreSQL and TimescaleDB have to build one
# Python tuple per row for `execute_values`, and InfluxDB one `Point` per row
# for line protocol. That construction is what using those interfaces costs,
# so it stays inside the timed region, and each of the three also measures it
# on its own and prints the share; the table above the chart collects them.
# Those three numbers are the *explicit* row construction and nothing more.
# Every client converts and serializes somewhere - `to_sql` still prepares and
# binds each row, `insert_df` still serializes the frame - and none of that is
# separable from outside, so a bar without a measured share is not a bar that
# is all engine. Read every write bar as "what it costs to get this panel into
# this engine through its Python client", never as the engine's ingest rate.
# - **Durability inside the timed region.** PostgreSQL and TimescaleDB commit
# synchronously, so their flush cost is inside the timed call. QuestDB (ILP
# + WAL) and InfluxDB acknowledge before the rows are queryable, so they
# poll to first-queryable via `wait_until_rows_visible` *inside* the timed
# call. Every write time therefore ends at the same event: the data is
# durable and readable.
# - **Reads, aggregations, and joins** (`time_read`) - mean of `TIMING_RUNS`
# runs after one untimed warm-up. **These are warm numbers**, on both the OS
# page cache and each server's buffer pool, and they flatter every engine
# that caches. A client cannot drop a server's cache, so a "cold" read here
# would be cold for the embedded engines and warm for the servers - neither
# cold nor comparable. Warm and uniform is the measurable choice; §2.4 says
# so where it reports these figures.
# - **Every read is validated** against the expected row count, exactly. Each
# query below is deterministic and its answer is known before it runs, so a
# near miss is a wrong answer rather than a tolerable one - and a query that
# silently returns a truncated result would otherwise post the fastest time
# on the chart.
#
# ## What the size column does and does not mean
#
# Unlike the timings, on-disk size is **not** measured the same way for every
# engine, because each engine only exposes its own accounting. Read the size
# column as "what this engine reports it is using", not as a like-for-like
# compression ratio:
#
# | Engine | Reported by | Counts |
# |--------|-------------|--------|
# | SQLite / DuckDB | file size on disk | the database file |
# | ClickHouse | `system.parts.bytes_on_disk` | **compressed** data parts |
# | QuestDB | `table_storage().diskSize` | the table's full on-disk footprint |
# | PostgreSQL / TimescaleDB | `pg_total_relation_size` / `hypertable_size` | table **plus indexes** (and, for the hypertable, all its chunks) |
# | InfluxDB | *not exposed* | client API reports no per-bucket size, so it is left empty rather than recorded as zero |
#
# Comparing ClickHouse's compressed parts against PostgreSQL's table-plus-index
# total is not a compression comparison. The file-format notebook
# (`20_storage_benchmark_file`) is where size is measured identically across
# contenders and can be compared directly.
#
# ## Quick Start
#
# ```bash
# # Embedded engines only. DuckDB and SQLite need nothing beyond the base install.
# uv run python 02_financial_data_universe/21_storage_benchmark_database.py
#
# # With the server engines as well, from the host (compose publishes their ports)
# docker compose --profile databases up -d timescaledb postgres clickhouse questdb influxdb
# uv run --extra db-benchmark python 02_financial_data_universe/21_storage_benchmark_database.py
# ```
#
# The `db-benchmark` extra also carries ArcticDB, which publishes no Linux ARM64 wheel
# and no source distribution. On ARM64 the install fails before the notebook can run,
# so the optional-import guard never gets the chance to skip it: use the Docker path
# above, whose `benchmark` image ships the server clients without ArcticDB.
#
# The scale is the `BENCHMARK_SCALE` parameter below, not an environment variable:
# the cell after the parameters cell writes the parameter back into the environment
# for `utils.storage_benchmarks` to read, so a `BENCHMARK_SCALE=...` prefix on the
# command line is overwritten before anything reads it. Change the scale by editing
# the parameter, or by injecting it with Papermill as CI does.
#
# Whichever engines answer, the coverage report near the bottom names the ones that
# produced the numbers on this page, and the ones that did not.
# %% [markdown]
# ## Setup
# %%
"""Storage Benchmark — Database engine comparison for financial time-series."""
import contextlib
import gc
import json
import os
import shutil
import sqlite3
import subprocess
import time as time_module
import urllib.parse
import urllib.request
from datetime import UTC, timedelta
from pathlib import Path
# %% [markdown]
# ### Declared parameters
#
# `BENCHMARK_SCALE` selects the panel size; the production setting is the one §2.4
# quotes, and CI overrides it to the small scale through Papermill. Every cell that
# reports a number also reports the scale that produced it, so a figure lifted off
# this page carries its own provenance.
#
# `RANGE_QUERY_SHARE` is the fraction of the panel the range query should select.
# It is a share rather than a fixed number of days because the panel's calendar span
# grows with the scale: a fixed seven-day window selects a fifth of the large panel
# and *all* of the small one, which would silently turn the range-query panel of the
# chart into a second copy of the full scan.
#
# `LOG_AXIS_RATIO` is the spread at which a chart panel switches to a logarithmic
# x-axis, measured on the bars actually drawn rather than assumed from a past run.
# %% tags=["parameters"]
BENCHMARK_SCALE = "L"
RANGE_QUERY_SHARE = 0.2
LOG_AXIS_RATIO = 10.0
# %% [markdown]
# `utils.storage_benchmarks` reads the scale from the environment when it is imported,
# so the variable has to be set before the import rather than passed to a function
# afterwards.
# %%
os.environ["BENCHMARK_SCALE"] = BENCHMARK_SCALE
import pandas as pd
import plotly.graph_objects as go
import polars as pl
from plotly.subplots import make_subplots
from utils.paths import display_path, get_output_dir
from utils.storage_benchmarks import (
ACTIVE_SCALE,
BENCHMARK_DIR,
DB_CONFIG,
N_ROWS_PER_SYMBOL,
N_SYMBOLS,
N_TICKS_QUOTES,
N_TICKS_TRADES,
TIMING_RUNS,
WAL_FLUSH_TIMEOUT,
BenchmarkResult,
estimate_memory_mb,
force_materialize_pandas,
force_materialize_polars,
generate_ohlcv_data,
generate_tick_data,
get_scale_config,
save_benchmark_results,
save_chart,
time_read,
time_write,
validate_result,
wait_until_rows_visible,
)
from utils.style import COLORS, show_plotly_with_alt
# %%
OUTPUT_DIR = get_output_dir(2, "storage_benchmark")
OUTPUT_DIR.mkdir(parents=True, exist_ok=True)
# %% [markdown]
# ## Check Available Databases
# %%
benchmark_status = {
# Embedded (always available if package installed)
"SQLite": {"expected": True, "tested": False, "category": "embedded"},
# DuckDB is imported unconditionally below, so its absence raises rather than
# lowering the bar this run is measured against.
"DuckDB": {"expected": True, "tested": False, "category": "embedded"},
"ArcticDB": {"expected": False, "tested": False, "category": "embedded"},
# Servers (need Docker)
"ClickHouse": {"expected": True, "tested": False, "category": "server"},
"QuestDB": {"expected": True, "tested": False, "category": "server"},
"TimescaleDB": {"expected": True, "tested": False, "category": "server"},
"InfluxDB": {"expected": True, "tested": False, "category": "server"},
"PostgreSQL": {"expected": True, "tested": False, "category": "server"},
"kdb+/PyKX": {"expected": False, "tested": False, "category": "hft"},
}
print("=" * 70)
print(f"DATABASE BENCHMARK - Scale: {ACTIVE_SCALE}")
print("=" * 70)
# %% [markdown]
# ### Which engines are present
#
# DuckDB ships in both the `benchmark` (ARM64) and `benchmark-full` (x86) images, and
# SQLite comes with Python, so both are always available and the notebook refuses to
# continue without DuckDB rather than quietly dropping it. ArcticDB is x86-only and
# lives in `benchmark-full`; where it is absent it is reported as absent.
# %%
try:
import duckdb # noqa: F401
except ImportError as exc:
raise ImportError(
"DuckDB is not available in the current image.\n"
"This notebook runs in the `benchmark` image:\n"
" docker compose --profile benchmark run --rm benchmark \\\n"
" python 02_financial_data_universe/21_storage_benchmark_database.py"
) from exc
print("\n### Embedded Databases")
print("[OK] DuckDB: Available")
print("[OK] SQLite: Available (built-in)")
# %% [markdown]
# The ArcticDB guard below catches `Exception`, not `ImportError`. An ArcticDB that is
# installed but cannot run raises something else entirely - `NotImplementedError`
# against an unsupported protobuf, for one - and a guard that catches only
# `ImportError` lets that end the notebook. An optional engine has to be optional in
# every way it can fail to be usable, not only in the one way it can be absent.
# %%
try:
import arcticdb as adb
HAS_ARCTICDB = True
benchmark_status["ArcticDB"]["expected"] = True
print("[OK] ArcticDB: Available")
except Exception as exc:
HAS_ARCTICDB = False
print(f"○ ArcticDB: unavailable, skipping ({type(exc).__name__}: {exc})")
# %%
# === Check Server Databases ===
print("\n### Server Databases (Docker required)")
# ClickHouse
try:
import clickhouse_connect
ch_client = clickhouse_connect.get_client(
host=DB_CONFIG["clickhouse"]["host"], port=DB_CONFIG["clickhouse"]["port"]
)
ch_client.query("SELECT 1")
HAS_CLICKHOUSE = True
print("[OK] ClickHouse: Available")
except Exception:
HAS_CLICKHOUSE = False
print("[FAIL] ClickHouse: Not available (start Docker)")
# QuestDB
try:
urllib.request.urlopen(
f"http://{DB_CONFIG['questdb']['host']}:{DB_CONFIG['questdb']['http_port']}/exec?query=SELECT%201",
timeout=2,
)
HAS_QUESTDB = True
print("[OK] QuestDB: Available")
except Exception:
HAS_QUESTDB = False
print("[FAIL] QuestDB: Not available (start Docker)")
# %%
# Check TimescaleDB and InfluxDB availability
try:
import psycopg2
ts_conn = psycopg2.connect(
host=DB_CONFIG["timescaledb"]["host"],
port=DB_CONFIG["timescaledb"]["port"],
user=DB_CONFIG["timescaledb"]["user"],
password=DB_CONFIG["timescaledb"]["password"],
database=DB_CONFIG["timescaledb"]["database"],
connect_timeout=2,
)
ts_conn.close()
HAS_TIMESCALEDB = True
print("[OK] TimescaleDB: Available")
except Exception:
HAS_TIMESCALEDB = False
print("[FAIL] TimescaleDB: Not available (start Docker)")
# InfluxDB
try:
from influxdb_client import InfluxDBClient
influx_test = InfluxDBClient(
url=f"http://{DB_CONFIG['influxdb']['host']}:{DB_CONFIG['influxdb']['port']}",
token=DB_CONFIG["influxdb"]["token"],
org=DB_CONFIG["influxdb"]["org"],
timeout=2000,
)
HAS_INFLUXDB = bool(influx_test.ready())
influx_test.close()
del influx_test
if HAS_INFLUXDB:
print("[OK] InfluxDB: Available")
else:
print("[FAIL] InfluxDB: Not ready")
except Exception:
HAS_INFLUXDB = False
print("[FAIL] InfluxDB: Not available (start Docker)")
# %%
# PostgreSQL (vanilla, separate from TimescaleDB)
try:
import psycopg2
pg_conn_check = psycopg2.connect(
host=DB_CONFIG["postgres"]["host"],
port=DB_CONFIG["postgres"]["port"],
user=DB_CONFIG["postgres"]["user"],
password=DB_CONFIG["postgres"]["password"],
database=DB_CONFIG["postgres"]["database"],
connect_timeout=2,
)
pg_conn_check.close()
HAS_POSTGRES = True
print("[OK] PostgreSQL: Available")
except Exception:
HAS_POSTGRES = False
print("[FAIL] PostgreSQL: Not available (start Docker)")
# %% [markdown]
# ### kdb+ via PyKX, in IPC mode
#
# The q binary and the licence are looked for *before* PyKX is imported, and the
# order is load-bearing. PyKX must never be imported in unlicensed mode
# (`PYKX_UNLICENSED=1`). An unlicensed import leaves the process in a state where
# later numpy work segfaults, and the crash lands in `generate_ohlcv_data`, nowhere
# near the import that caused it, so it reads as a data bug. Without a licence there
# is nothing to benchmark, so PyKX is not imported at all and kdb+ is skipped.
# %%
HAS_PYKX = False
Q_BINARY: Path | None = None
# Check multiple locations for q binary (host install or Docker mount)
Q_BINARY_LOCATIONS = [
Path.home() / ".kx" / "bin" / "q", # Standard location
Path("/opt/kx/bin/q"), # Alternative system location
]
KX_LICENSE_DIRS = [Path.home() / ".pykx", Path.home() / ".kx"]
for q_path in Q_BINARY_LOCATIONS:
if q_path.exists() and q_path.is_file():
Q_BINARY = q_path
break
KX_LICENSE_FILE = next(
(d / f for d in KX_LICENSE_DIRS for f in ("kc.lic", "k4.lic") if (d / f).exists()), None
)
if Q_BINARY is None or KX_LICENSE_FILE is None:
print("○ PyKX/kdb+: Optional (not configured) — skipping, no result will be claimed")
if Q_BINARY is None:
print(" → Get a free personal edition: https://kx.com/kdb-personal-edition-download/")
print(f" → Install the q binary to: {display_path(Q_BINARY_LOCATIONS[0])}")
if KX_LICENSE_FILE is None:
print(f" → Place the licence file (kc.lic) in: {display_path(KX_LICENSE_DIRS[0])}")
else:
try:
# q resolves its licence from QLIC, not from the file merely existing on
# disk: without this the q process starts and dies "no license loaded".
os.environ["QLIC"] = str(KX_LICENSE_FILE.parent)
import pykx as kx
HAS_PYKX = True
benchmark_status["kdb+/PyKX"]["expected"] = True
print(
f"[OK] PyKX/kdb+ {kx.__version__}: Available "
f"(IPC mode, licence {display_path(KX_LICENSE_FILE)})"
)
except Exception as exc:
# As with ArcticDB: an optional engine that is present but unusable must not
# end the run, whatever it raises on the way out.
print(f"○ PyKX/kdb+: unavailable, skipping ({type(exc).__name__}: {exc})")
# %% [markdown]
# `expected` marks an engine as one the chapter compares. `available` says it answered
# the availability check on this machine. They are different questions, and conflating
# them turns a laptop with no Docker running into a run that reports five failures.
# The three counts below, and the coverage report at the end, are read off the status
# table rather than written as literals, so adding an engine to that table is most of
# what adding an engine takes.
# %%
AVAILABILITY = {
"SQLite": True,
"DuckDB": True,
"ArcticDB": HAS_ARCTICDB,
"ClickHouse": HAS_CLICKHOUSE,
"QuestDB": HAS_QUESTDB,
"TimescaleDB": HAS_TIMESCALEDB,
"InfluxDB": HAS_INFLUXDB,
"PostgreSQL": HAS_POSTGRES,
"kdb+/PyKX": HAS_PYKX,
}
assert set(AVAILABILITY) == set(benchmark_status), (
"every engine in the status table needs an availability flag"
)
for _engine, _is_available in AVAILABILITY.items():
benchmark_status[_engine]["available"] = _is_available
n_embedded = sum(v["available"] for v in benchmark_status.values() if v["category"] == "embedded")
n_servers = sum(v["available"] for v in benchmark_status.values() if v["category"] == "server")
n_hft = sum(v["available"] for v in benchmark_status.values() if v["category"] == "hft")
print(f"\n[OK] {n_embedded} embedded + {n_servers} server + {n_hft} HFT database(s) available")
if n_servers == 0:
print("\n[WARN] No server databases. Start Docker containers:")
print(
" docker compose --profile benchmark up -d timescaledb clickhouse questdb influxdb postgres"
)
# %% [markdown]
# ## Generate Test Data
# %%
scale_cfg = get_scale_config(ACTIVE_SCALE)
print(f"\nTarget: {scale_cfg['target_memory']} in-memory")
print(f"OHLCV: {N_SYMBOLS} symbols × {N_ROWS_PER_SYMBOL:,} rows/symbol")
print(f"Tick: {N_TICKS_TRADES:,} trades, {N_TICKS_QUOTES:,} quotes")
print("\n=== Generating synthetic data ===\n")
# Generate OHLCV data
ohlcv_df = generate_ohlcv_data(n_symbols=N_SYMBOLS, n_rows=N_ROWS_PER_SYMBOL)
total_rows = len(ohlcv_df)
memory_mb = estimate_memory_mb(ohlcv_df)
print(f"OHLCV data: {total_rows:,} rows ({memory_mb:.2f} MB)")
# Generate tick data for ASOF joins
trades_df, quotes_df = generate_tick_data(
n_trades=N_TICKS_TRADES, n_quotes=N_TICKS_QUOTES, n_symbols=N_SYMBOLS
)
print(f"Tick data: {len(trades_df):,} trades, {len(quotes_df):,} quotes")
# pandas versions (some tools require pandas)
ohlcv_pandas = ohlcv_df.to_pandas()
trades_pandas = trades_df.to_pandas()
quotes_pandas = quotes_df.to_pandas()
# Results collection
results: list[BenchmarkResult] = []
# %% [markdown]
# ### Measuring the row construction that is written out
#
# PostgreSQL, TimescaleDB and InfluxDB reach their servers through interfaces that
# take one Python object per row, so building those objects is part of what writing
# through them costs and it stays inside the timed write. Because that construction
# is a separate expression, it can also be timed on its own, which the helper below
# does. What comes back is a lower bound on each of those three clients' share, not
# the whole of it: the client still converts and serializes what it is handed, and
# so does every client whose interface takes a frame. Nothing here measures that.
# %%
client_side_build: dict[str, float] = {}
def measure_row_build(engine: str, build) -> float:
"""Time an engine's client-side row construction in isolation."""
gc.collect()
started = time_module.perf_counter()
rows = build()
elapsed = time_module.perf_counter() - started
del rows
gc.collect()
client_side_build[engine] = elapsed
return elapsed
def report_row_build(engine: str, write_seconds: float) -> None:
"""Print an engine's client-side share of its own write time."""
build_seconds = client_side_build[engine]
print(
f" Client-side row building: {build_seconds:.3f}s of the {write_seconds:.3f}s "
f"write ({build_seconds / write_seconds:.0%})"
)
# %% [markdown]
# ## Benchmark Windows
#
# Two derived quantities that every engine below reuses, so that all engines answer
# the *same* question and can be validated against the same expected row count:
#
# - **Range query**: the leading sessions of the panel, the slice a backtest walks.
# The window is sized as a share of the panel rather than as a fixed number of
# calendar days, because the panel's span is a function of the scale. Seven
# calendar days is a fifth of the large panel but covers the small panel entirely,
# and a range query that selects every row is a full scan wearing a different
# name: the chart would show two panels measuring one query and the reader would
# have no way to tell. Sizing by share keeps the question the same at every scale,
# and the cell below prints the share it actually achieved.
# - **Aggregation**: minute bars resampled to **daily** bars. Every engine buckets
# by day; bucketing minute data by minute would be a near-identity for the engines
# that did it and a real reduction for the ones that did not, which is not a
# comparison.
# %%
session_first_ts = (
ohlcv_df.group_by(pl.col("timestamp").dt.date().alias("session"))
.agg(pl.col("timestamp").min().alias("first_ts"))
.sort("session")["first_ts"]
.to_list()
)
n_sessions = len(session_first_ts)
range_sessions = max(1, round(RANGE_QUERY_SHARE * n_sessions))
range_start = ohlcv_df["timestamp"].min()
if range_sessions < n_sessions:
range_end = session_first_ts[range_sessions]
else:
range_end = ohlcv_df["timestamp"].max() + timedelta(minutes=1)
range_expected_rows = ohlcv_df.filter(
(pl.col("timestamp") >= range_start) & (pl.col("timestamp") < range_end)
).height
range_share = range_expected_rows / total_rows
agg_expected_rows = ohlcv_df.select(
pl.struct("symbol", pl.col("timestamp").dt.truncate("1d")).n_unique()
).item()
print(f"Panel spans {n_sessions} trading sessions: {range_start} → {ohlcv_df['timestamp'].max()}")
print(
f"Range query : first {range_sessions} of {n_sessions} sessions, "
f"{range_start} ≤ timestamp < {range_end} "
f"→ {range_expected_rows:,} rows ({range_share:.1%} of panel)"
)
print(f"Aggregation : minute bars → {agg_expected_rows:,} daily bars")
if range_sessions >= n_sessions:
print(
" [WARN] The panel is only one window wide at this scale, so the range query "
"reads every row and its chart panel repeats the full scan."
)
# %% [markdown]
# ---
# # Part 1: Embedded Databases
# ---
# %% [markdown]
# ## SQLite (Embedded RDBMS)
#
# SQLite is an embedded relational database:
# - ACID compliant, single-file database
# - Good for transactional workloads (OLTP)
# - Limited analytical query optimization
# - No native ASOF join
# %%
print("\n" + "=" * 70)
print("SQLITE BENCHMARK")
print("=" * 70)
benchmark_status["SQLite"]["tested"] = True
sqlite_path = BENCHMARK_DIR / f"ohlcv_{ACTIVE_SCALE.lower()}.db"
# Write
def write_sqlite():
if sqlite_path.exists():
sqlite_path.unlink()
with contextlib.closing(sqlite3.connect(sqlite_path)) as conn:
ohlcv_pandas.to_sql("ohlcv", conn, if_exists="replace", index=False)
conn.execute("CREATE INDEX IF NOT EXISTS idx_symbol_timestamp ON ohlcv(symbol, timestamp)")
write_time, _ = time_write(write_sqlite)
sqlite_size = sqlite_path.stat().st_size
results.append(BenchmarkResult("SQLite", "write", write_time, sqlite_size, total_rows))
# %% [markdown]
# ### SQLite Read and Aggregation
# %%
def read_sqlite():
with contextlib.closing(sqlite3.connect(sqlite_path)) as conn:
df = pd.read_sql("SELECT * FROM ohlcv", conn)
return force_materialize_pandas(df)
read_time, sqlite_result = time_read(read_sqlite)
validate_result(sqlite_result, total_rows, "SQLite read")
results.append(BenchmarkResult("SQLite", "read", read_time, sqlite_size, total_rows))
# %% [markdown]
# ### SQLite Range Query
#
# SQLite stores the timestamps as ISO-8601 text, which sorts lexicographically,
# so the string bounds below use the covering `(symbol, timestamp)` index.
# %%
SQLITE_TS_FMT = "%Y-%m-%d %H:%M:%S"
def sqlite_range_query():
with contextlib.closing(sqlite3.connect(sqlite_path)) as conn:
df = pd.read_sql(
"SELECT * FROM ohlcv WHERE timestamp >= ? AND timestamp < ?",
conn,
params=(range_start.strftime(SQLITE_TS_FMT), range_end.strftime(SQLITE_TS_FMT)),
)
return force_materialize_pandas(df)
range_time, sqlite_range = time_read(sqlite_range_query)
validate_result(sqlite_range, range_expected_rows, "SQLite range query")
results.append(BenchmarkResult("SQLite", "range_query", range_time, 0, len(sqlite_range)))
# %% [markdown]
# ### SQLite OHLCV Aggregation
#
# `open` is the first price of the session and `close` the last, ordered by
# time — not `MIN(open)` / `MAX(close)`, which are a different (and wrong)
# statistic. SQLite has no `first`/`last` aggregate, so the OHLCV reduction
# needs window functions. That costs SQLite time relative to the engines with
# native `first`/`last`, and that cost is the honest answer to "what does this
# aggregation take on SQLite?"
# %%
def sqlite_aggregation():
with contextlib.closing(sqlite3.connect(sqlite_path)) as conn:
query = """
SELECT DISTINCT symbol, date(timestamp) as bar_date,
FIRST_VALUE(open) OVER w as open,
MAX(high) OVER w as high,
MIN(low) OVER w as low,
LAST_VALUE(close) OVER w as close,
SUM(volume) OVER w as volume
FROM ohlcv
WINDOW w AS (
PARTITION BY symbol, date(timestamp) ORDER BY timestamp
ROWS BETWEEN UNBOUNDED PRECEDING AND UNBOUNDED FOLLOWING
)
"""
return pd.read_sql(query, conn)
agg_time, agg_result = time_read(sqlite_aggregation)
validate_result(agg_result, agg_expected_rows, "SQLite aggregation")
results.append(BenchmarkResult("SQLite", "aggregation", agg_time, 0, len(agg_result)))
print(f"\nSQLite: {sqlite_size / 1e6:.1f} MB")
print(f" Write: {write_time:.3f}s ({total_rows / write_time / 1e6:.2f}M rows/s)")
print(f" Read: {read_time:.3f}s ({total_rows / read_time / 1e6:.2f}M rows/s)")
print(f" Range query: {range_time:.3f}s ({len(sqlite_range):,} rows)")
print(f" Aggregation: {agg_time:.3f}s ({len(agg_result):,} daily bars)")
print(" Note: No native ASOF join")
# %% [markdown]
# ## DuckDB (Embedded Analytics)
#
# DuckDB is designed for analytical workloads (OLAP):
# - Columnar storage, vectorized execution
# - Native ASOF join support (v1.1+)
# - Zero-copy reads from Parquet
# - Out-of-core processing for data larger than RAM
# %%
print("\n" + "=" * 70)
print("DUCKDB BENCHMARK")
print("=" * 70)
benchmark_status["DuckDB"]["tested"] = True
duckdb_path = BENCHMARK_DIR / f"ohlcv_{ACTIVE_SCALE.lower()}.duckdb"
parquet_path = BENCHMARK_DIR / f"ohlcv_{ACTIVE_SCALE.lower()}.parquet"
# Save to Parquet for DuckDB's preferred workflow
ohlcv_df.write_parquet(parquet_path)
# Write DuckDB native
def write_duckdb():
if duckdb_path.exists():
duckdb_path.unlink()
conn = duckdb.connect(str(duckdb_path))
conn.execute("CREATE TABLE ohlcv AS SELECT * FROM read_parquet(?)", [str(parquet_path)])
conn.close()
write_time, _ = time_write(write_duckdb)
duckdb_size = duckdb_path.stat().st_size
results.append(BenchmarkResult("DuckDB", "write", write_time, duckdb_size, total_rows))
# Read
def read_duckdb():
conn = duckdb.connect(str(duckdb_path), read_only=True)
df = conn.execute("SELECT * FROM ohlcv").pl()
conn.close()
return force_materialize_polars(df)
read_time, duckdb_result = time_read(read_duckdb)
validate_result(duckdb_result, total_rows, "DuckDB read")
results.append(BenchmarkResult("DuckDB", "read", read_time, duckdb_size, total_rows))
# %%
# Range query
def duckdb_range_query():
conn = duckdb.connect(str(duckdb_path), read_only=True)
df = conn.execute(
"SELECT * FROM ohlcv WHERE timestamp >= ? AND timestamp < ?", [range_start, range_end]
).pl()
conn.close()
return force_materialize_polars(df)
range_time, duckdb_range = time_read(duckdb_range_query)
validate_result(duckdb_range, range_expected_rows, "DuckDB range query")
results.append(BenchmarkResult("DuckDB", "range_query", range_time, 0, len(duckdb_range)))
# %%
# Aggregation
def duckdb_aggregation():
conn = duckdb.connect(str(duckdb_path), read_only=True)
result = conn.execute("""
SELECT symbol, date_trunc('day', timestamp) as bar_date,
FIRST(open ORDER BY timestamp) as open, MAX(high) as high,
MIN(low) as low, LAST(close ORDER BY timestamp) as close,
SUM(volume) as volume
FROM ohlcv GROUP BY symbol, bar_date ORDER BY symbol, bar_date
""").pl()
conn.close()
return result
agg_time, agg_result = time_read(duckdb_aggregation)
validate_result(agg_result, agg_expected_rows, "DuckDB aggregation")
results.append(BenchmarkResult("DuckDB", "aggregation", agg_time, 0, len(agg_result)))
# ASOF Join
trades_path = BENCHMARK_DIR / f"trades_{ACTIVE_SCALE.lower()}.parquet"
quotes_path = BENCHMARK_DIR / f"quotes_{ACTIVE_SCALE.lower()}.parquet"
trades_df.sort(["symbol", "timestamp"]).write_parquet(trades_path)
quotes_df.sort(["symbol", "timestamp"]).write_parquet(quotes_path)
def duckdb_asof():
conn = duckdb.connect()
result = conn.execute(f"""
SELECT t.*, q.bid, q.ask, q.bid_size, q.ask_size
FROM read_parquet('{trades_path}') t
ASOF LEFT JOIN read_parquet('{quotes_path}') q
ON t.symbol = q.symbol AND t.timestamp >= q.timestamp
""").pl()
conn.close()
return result
asof_time, asof_result = time_read(duckdb_asof)
validate_result(asof_result, N_TICKS_TRADES, "DuckDB ASOF join")
results.append(BenchmarkResult("DuckDB", "asof_join", asof_time, 0, len(asof_result)))
print(f"\nDuckDB: {duckdb_size / 1e6:.1f} MB")
print(f" Write: {write_time:.3f}s ({total_rows / write_time / 1e6:.2f}M rows/s)")
print(f" Read: {read_time:.3f}s ({total_rows / read_time / 1e6:.2f}M rows/s)")
print(f" Range query: {range_time:.3f}s ({len(duckdb_range):,} rows)")
print(f" Aggregation: {agg_time:.3f}s ({len(agg_result):,} daily bars)")
print(f" ASOF Join: {asof_time:.3f}s ({N_TICKS_TRADES / asof_time / 1e6:.2f}M trades/s)")
# %% [markdown]
# ## ArcticDB (Versioned DataFrames)
#
# ArcticDB is designed for versioned time-series storage:
# - "Git for DataFrames" - version history, time travel
# - Optimized for financial time-series
# - LMDB backend (local), S3/Azure (cloud)
#
# ArcticDB answers the read and the aggregation but not the range query. Its
# server-side `date_range` filter keys off a datetime *index*, and this panel carries
# `timestamp` as a column, which is the canonical schema; pushing the filter down
# would need a different write layout from the one timed above. So ArcticDB is absent
# from the range-query panel rather than being timed on a client-side filter no other
# engine pays.
# %%
if HAS_ARCTICDB:
print("\n" + "=" * 70)
print("ARCTICDB BENCHMARK")
print("=" * 70)
benchmark_status["ArcticDB"]["tested"] = True
arctic_path = BENCHMARK_DIR / f"arctic_{ACTIVE_SCALE.lower()}"
if arctic_path.exists():
shutil.rmtree(arctic_path)
ac = adb.Arctic(f"lmdb://{arctic_path}")
lib = ac.get_library("benchmark", create_if_missing=True)
def write_arctic():
lib.write("ohlcv", ohlcv_pandas, prune_previous_versions=True)
write_time, _ = time_write(write_arctic)
arctic_size = sum(f.stat().st_size for f in arctic_path.rglob("*") if f.is_file())
results.append(BenchmarkResult("ArcticDB", "write", write_time, arctic_size, total_rows))
def read_arctic():
df = lib.read("ohlcv").data
return force_materialize_pandas(df)
read_time, arctic_result = time_read(read_arctic, n_runs=min(2, TIMING_RUNS))
validate_result(arctic_result, total_rows, "ArcticDB read")
results.append(BenchmarkResult("ArcticDB", "read", read_time, arctic_size, total_rows))
def arctic_aggregation():
df = lib.read("ohlcv").data
return (
df.groupby(["symbol", pd.Grouper(key="timestamp", freq="D")])
.agg(
open=("open", "first"),
high=("high", "max"),
low=("low", "min"),
close=("close", "last"),
volume=("volume", "sum"),
)
.reset_index()
)
agg_time, agg_result = time_read(arctic_aggregation, n_runs=min(2, TIMING_RUNS))
validate_result(agg_result, agg_expected_rows, "ArcticDB aggregation")
results.append(BenchmarkResult("ArcticDB", "aggregation", agg_time, 0, len(agg_result)))
print(f"\nArcticDB: {arctic_size / 1e6:.1f} MB")
print(f" Write: {write_time:.3f}s ({total_rows / write_time / 1e6:.2f}M rows/s)")
print(f" Read: {read_time:.3f}s ({total_rows / read_time / 1e6:.2f}M rows/s)")
print(f" Aggregation: {agg_time:.3f}s ({len(agg_result):,} daily bars)")
print(" Note: Versioning enabled (time travel supported)")
shutil.rmtree(arctic_path)
else:
print("\nArcticDB benchmark skipped — install via the benchmark-full image (x86 only).")
# %% [markdown]
# ---
# # Part 2: Server Databases (Docker Required)
# ---
# %% [markdown]
# ## ClickHouse (OLAP Analytics)
#
# ClickHouse excels at:
# - Massive aggregations (billions of rows/second)
# - Columnar compression (10-15x)
# - Native ASOF JOIN
# %%
if HAS_CLICKHOUSE:
print("\n" + "=" * 70)
print("CLICKHOUSE BENCHMARK")
print("=" * 70)
benchmark_status["ClickHouse"]["tested"] = True
ch_client.command("DROP TABLE IF EXISTS ohlcv_benchmark")
ch_client.command("""
CREATE TABLE ohlcv_benchmark (
timestamp DateTime64(3), symbol String,
open Float64, high Float64, low Float64, close Float64,
volume Int64, vwap Float64, num_trades Int32
) ENGINE = MergeTree()
PARTITION BY toYYYYMM(timestamp) ORDER BY (symbol, timestamp)
""")
# Write
def write_clickhouse():
ch_client.insert_df("ohlcv_benchmark", ohlcv_pandas)
ch_write_time, _ = time_write(write_clickhouse)
ch_size = (
ch_client.query(
"SELECT sum(bytes_on_disk) FROM system.parts WHERE table = 'ohlcv_benchmark'"
).result_set[0][0]
or 0
)
results.append(BenchmarkResult("ClickHouse", "write", ch_write_time, ch_size, total_rows))
# Read
def read_clickhouse():
return ch_client.query_df("SELECT * FROM ohlcv_benchmark")
ch_read_time, ch_result = time_read(read_clickhouse, n_runs=min(3, TIMING_RUNS))
validate_result(ch_result, total_rows, "ClickHouse read")
results.append(BenchmarkResult("ClickHouse", "read", ch_read_time, ch_size, total_rows))
# %%
if HAS_CLICKHOUSE:
# Range query
def clickhouse_range_query():
return ch_client.query_df(
"SELECT * FROM ohlcv_benchmark WHERE timestamp >= %(start)s AND timestamp < %(end)s",
parameters={"start": range_start, "end": range_end},
)
ch_range_time, ch_range_result = time_read(clickhouse_range_query, n_runs=min(3, TIMING_RUNS))
validate_result(ch_range_result, range_expected_rows, "ClickHouse range query")
results.append(
BenchmarkResult("ClickHouse", "range_query", ch_range_time, 0, len(ch_range_result))
)
# %%
if HAS_CLICKHOUSE:
# Aggregation — bucket to DAY, matching every other engine.
def clickhouse_ohlcv():
return ch_client.query_df("""
SELECT symbol, toStartOfDay(timestamp) as bar_time,
argMin(open, timestamp) as open, max(high) as high,
min(low) as low, argMax(close, timestamp) as close,
sum(volume) as volume
FROM ohlcv_benchmark GROUP BY symbol, bar_time ORDER BY symbol, bar_time
""")
ch_agg_time, ch_agg_result = time_read(clickhouse_ohlcv, n_runs=min(3, TIMING_RUNS))
validate_result(ch_agg_result, agg_expected_rows, "ClickHouse aggregation")
results.append(BenchmarkResult("ClickHouse", "aggregation", ch_agg_time, 0, len(ch_agg_result)))
# %%
if HAS_CLICKHOUSE:
# ASOF JOIN
ch_client.command("DROP TABLE IF EXISTS ch_trades")
ch_client.command("DROP TABLE IF EXISTS ch_quotes")
ch_client.command("""
CREATE TABLE ch_trades (timestamp DateTime64(9), symbol String, price Float64, size Int64)
ENGINE = MergeTree() ORDER BY (symbol, timestamp)
""")
ch_client.command("""
CREATE TABLE ch_quotes (timestamp DateTime64(9), symbol String, bid Float64, ask Float64, bid_size Int64, ask_size Int64)
ENGINE = MergeTree() ORDER BY (symbol, timestamp)
""")
trades_sorted_pd = (
trades_df.sort(["symbol", "timestamp"])
.select(["timestamp", "symbol", "price", "size"])
.to_pandas()
)
quotes_sorted_pd = (
quotes_df.sort(["symbol", "timestamp"])
.select(["timestamp", "symbol", "bid", "ask", "bid_size", "ask_size"])
.to_pandas()
)
ch_client.insert_df("ch_trades", trades_sorted_pd)
ch_client.insert_df("ch_quotes", quotes_sorted_pd)
# %%
if HAS_CLICKHOUSE:
def clickhouse_asof():
return ch_client.query_df("""
SELECT t.timestamp, t.symbol, t.price, t.size, q.bid, q.ask
FROM ch_trades t ASOF LEFT JOIN ch_quotes q
ON t.symbol = q.symbol AND t.timestamp >= q.timestamp
""")
ch_asof_time, ch_asof_result = time_read(clickhouse_asof, n_runs=min(3, TIMING_RUNS))
validate_result(ch_asof_result, N_TICKS_TRADES, "ClickHouse ASOF join")
results.append(BenchmarkResult("ClickHouse", "asof_join", ch_asof_time, 0, len(ch_asof_result)))
print(f"\nClickHouse: {ch_size / 1e6:.1f} MB")
print(f" Write: {ch_write_time:.3f}s | Read: {ch_read_time:.3f}s")
print(f" Range query: {ch_range_time:.3f}s ({len(ch_range_result):,} rows)")
print(f" Aggregation: {ch_agg_time:.3f}s ({len(ch_agg_result):,} daily bars)")
print(f" ASOF Join: {ch_asof_time:.3f}s ({len(ch_asof_result):,} rows)")
ch_client.command("DROP TABLE IF EXISTS ohlcv_benchmark")
ch_client.command("DROP TABLE IF EXISTS ch_trades")
ch_client.command("DROP TABLE IF EXISTS ch_quotes")
else:
print("\n⊘ ClickHouse benchmark skipped")
# %% [markdown]
# ## QuestDB (High-Throughput Time-Series)
#
# QuestDB is optimized for:
# - Ultra-high ingestion (1M+ rows/second via ILP)
# - Time-series specific SQL extensions (SAMPLE BY)
# - Native ASOF JOIN
# %% [markdown]
# QuestDB ingests over ILP, which acknowledges before the rows are queryable. The WAL
# commit is therefore polled *inside* the timed region, so QuestDB's write ends where
# PostgreSQL's does: when the data can be read back. A fixed sleep outside the timed
# call, which is what this used to do, charges QuestDB nothing for durability.
#
# The panel is handed to `Sender.dataframe` as a block, so QuestDB pays no per-row
# Python cost and does not appear in the client-side table.
# %%
if HAS_QUESTDB:
print("\n" + "=" * 70)
print("QUESTDB BENCHMARK")
print("=" * 70)
benchmark_status["QuestDB"]["tested"] = True
from questdb.ingress import Sender
def questdb_query(sql, limit: str | None = None):
"""Run SQL over QuestDB's HTTP endpoint.
`/exec` caps the JSON result set unless an explicit `limit` is given, so
every full-result query below passes one. Without it the endpoint
returns a truncated page: a fast time on a wrong answer.
"""
url = (
f"http://{DB_CONFIG['questdb']['host']}:{DB_CONFIG['questdb']['http_port']}"
f"/exec?query={urllib.parse.quote(sql)}"
)
if limit is not None:
url += f"&limit={limit}"
response = urllib.request.urlopen(url, timeout=600)
return json.loads(response.read())
def questdb_row_count() -> int:
try:
payload = questdb_query("SELECT count() FROM ohlcv_benchmark")
except Exception:
return 0 # table not created yet
dataset = payload.get("dataset") or [[0]]
return int(dataset[0][0])
with contextlib.suppress(Exception):
questdb_query("DROP TABLE IF EXISTS ohlcv_benchmark")
questdb_query("""
CREATE TABLE IF NOT EXISTS ohlcv_benchmark (
timestamp TIMESTAMP, symbol SYMBOL,
open DOUBLE, high DOUBLE, low DOUBLE, close DOUBLE,
volume LONG, vwap DOUBLE, num_trades INT
) timestamp(timestamp) PARTITION BY DAY WAL
""")
def write_questdb():
with Sender.from_conf(
f"http::addr={DB_CONFIG['questdb']['host']}:{DB_CONFIG['questdb']['http_port']};"
) as sender:
df_insert = ohlcv_pandas.copy()
df_insert["timestamp"] = pd.to_datetime(df_insert["timestamp"])
sender.dataframe(
df_insert, table_name="ohlcv_benchmark", symbols=["symbol"], at="timestamp"
)
return wait_until_rows_visible(
questdb_row_count, total_rows, timeout=max(60.0, WAL_FLUSH_TIMEOUT * 20)
)
qdb_write_time, qdb_visible = time_write(write_questdb)
assert qdb_visible == total_rows, f"QuestDB ingested {qdb_visible:,} of {total_rows:,}"
qdb_size = int(
questdb_query(
"SELECT sum(diskSize) FROM table_storage() WHERE tableName = 'ohlcv_benchmark'"
)["dataset"][0][0]
or 0
)
results.append(BenchmarkResult("QuestDB", "write", qdb_write_time, qdb_size, total_rows))
# %%
if HAS_QUESTDB:
# Read — explicit limit so the endpoint returns the whole table, then
# validated like every other engine's read.
def read_questdb():
# One past the expected count: asking for exactly the expected number would
# truncate an over-large result to exactly right and validate it as correct.
return questdb_query("SELECT * FROM ohlcv_benchmark", limit=f"1,{total_rows + 1}")
qdb_read_time, qdb_result = time_read(read_questdb, n_runs=min(3, TIMING_RUNS))
validate_result(qdb_result, total_rows, "QuestDB read")
results.append(BenchmarkResult("QuestDB", "read", qdb_read_time, qdb_size, total_rows))
# Range query
def questdb_range_query():
sql = (
"SELECT * FROM ohlcv_benchmark "
f"WHERE timestamp >= '{range_start.isoformat()}' "
f"AND timestamp < '{range_end.isoformat()}'"
)
return questdb_query(sql, limit=f"1,{range_expected_rows + 1}")
qdb_range_time, qdb_range_result = time_read(questdb_range_query, n_runs=min(3, TIMING_RUNS))
validate_result(qdb_range_result, range_expected_rows, "QuestDB range query")
results.append(
BenchmarkResult(
"QuestDB", "range_query", qdb_range_time, 0, len(qdb_range_result["dataset"])
)
)
# OHLCV aggregation (SAMPLE BY) — daily buckets, matching every other engine
def questdb_ohlcv():
return questdb_query(
"""
SELECT symbol, timestamp as bar_time,
first(open) as open, max(high) as high, min(low) as low,
last(close) as close, sum(volume) as volume
FROM ohlcv_benchmark SAMPLE BY 1d ALIGN TO CALENDAR
""",
limit=f"1,{agg_expected_rows + 1}",
)
qdb_agg_time, qdb_agg_result = time_read(questdb_ohlcv, n_runs=min(3, TIMING_RUNS))
validate_result(qdb_agg_result, agg_expected_rows, "QuestDB aggregation")
qdb_agg_rows = len(qdb_agg_result["dataset"])
results.append(BenchmarkResult("QuestDB", "aggregation", qdb_agg_time, 0, qdb_agg_rows))
print(f"\nQuestDB: {qdb_size / 1e6:.1f} MB")
print(f" Write (ILP, to queryable): {qdb_write_time:.3f}s | Read: {qdb_read_time:.3f}s")
print(f" Range query: {qdb_range_time:.3f}s ({len(qdb_range_result['dataset']):,} rows)")
print(f" OHLCV aggregation: {qdb_agg_time:.3f}s ({qdb_agg_rows:,} daily bars)")
with contextlib.suppress(Exception):
questdb_query("DROP TABLE IF EXISTS ohlcv_benchmark")
else:
print("\n⊘ QuestDB benchmark skipped")
# %% [markdown]
# ## TimescaleDB (PostgreSQL + Time-Series)
#
# TimescaleDB combines PostgreSQL with time-series optimizations:
# - Hypertables (automatic partitioning)
# - time_bucket() for aggregations
# - Full SQL + relational integrity
# %%
if HAS_TIMESCALEDB:
print("\n" + "=" * 70)
print("TIMESCALEDB BENCHMARK")
print("=" * 70)
benchmark_status["TimescaleDB"]["tested"] = True
from psycopg2.extras import execute_values
conn = psycopg2.connect(
host=DB_CONFIG["timescaledb"]["host"],
port=DB_CONFIG["timescaledb"]["port"],
user=DB_CONFIG["timescaledb"]["user"],
password=DB_CONFIG["timescaledb"]["password"],
database=DB_CONFIG["timescaledb"]["database"],
)
conn.autocommit = True
cur = conn.cursor()
cur.execute("CREATE EXTENSION IF NOT EXISTS timescaledb CASCADE;")
cur.execute("DROP TABLE IF EXISTS ohlcv_benchmark CASCADE;")
cur.execute("""
CREATE TABLE ohlcv_benchmark (
timestamp TIMESTAMPTZ NOT NULL, symbol TEXT NOT NULL,
open DOUBLE PRECISION, high DOUBLE PRECISION, low DOUBLE PRECISION,
close DOUBLE PRECISION, volume BIGINT, vwap DOUBLE PRECISION, num_trades INTEGER
);
""")
cur.execute(
"SELECT create_hypertable('ohlcv_benchmark', 'timestamp', chunk_time_interval => INTERVAL '1 day');"
)
# %% [markdown]
# ### TimescaleDB write
#
# `execute_values` takes a sequence of Python tuples, so the panel has to be turned
# into one tuple per row before any of it reaches the server. That construction uses
# `itertuples`, not `iterrows`: `iterrows` boxes each row as a Series first, which
# costs an order of magnitude more to produce tuples that compare equal. The
# difference is pandas rather than the database, and it was being charged to this
# engine's write bar. The cell below prints what the construction cost on this run.
#
# On size, `hypertable_size()` rather than `pg_total_relation_size()`: a hypertable's
# rows live in child chunks, so the parent relation is empty and
# `pg_total_relation_size` reports about 16 kB of catalog overhead instead of data.
# %%
if HAS_TIMESCALEDB:
def build_timescaledb_rows():
return [
(
r.timestamp,
r.symbol,
r.open,
r.high,
r.low,
r.close,
int(r.volume),
r.vwap,
int(r.num_trades),
)
for r in ohlcv_pandas.itertuples(index=False)
]
def write_timescaledb():
execute_values(
cur,
"""
INSERT INTO ohlcv_benchmark (timestamp, symbol, open, high, low, close, volume, vwap, num_trades) VALUES %s
""",
build_timescaledb_rows(),
)
ts_write_time, _ = time_write(write_timescaledb)
measure_row_build("TimescaleDB", build_timescaledb_rows)
cur.execute("SELECT hypertable_size('ohlcv_benchmark');")
ts_size = cur.fetchone()[0]
results.append(BenchmarkResult("TimescaleDB", "write", ts_write_time, ts_size, total_rows))
# %%
if HAS_TIMESCALEDB:
# Read
def read_timescaledb():
cur.execute("SELECT * FROM ohlcv_benchmark;")
return cur.fetchall()
ts_read_time, ts_result = time_read(read_timescaledb, n_runs=min(3, TIMING_RUNS))
validate_result(ts_result, total_rows, "TimescaleDB read")
results.append(BenchmarkResult("TimescaleDB", "read", ts_read_time, ts_size, total_rows))
# %%
if HAS_TIMESCALEDB:
# Range query
def timescaledb_range_query():
cur.execute(
"SELECT * FROM ohlcv_benchmark WHERE timestamp >= %s AND timestamp < %s;",
(range_start, range_end),
)
return cur.fetchall()
ts_range_time, ts_range_result = time_read(timescaledb_range_query, n_runs=min(3, TIMING_RUNS))
validate_result(ts_range_result, range_expected_rows, "TimescaleDB range query")
results.append(
BenchmarkResult("TimescaleDB", "range_query", ts_range_time, 0, len(ts_range_result))
)
# Aggregation (time_bucket) — daily buckets, matching every other engine
def timescaledb_ohlcv():
cur.execute("""
SELECT symbol, time_bucket('1 day', timestamp) as bar_time,
first(open, timestamp) as open, max(high) as high,
min(low) as low, last(close, timestamp) as close, sum(volume) as volume
FROM ohlcv_benchmark GROUP BY symbol, bar_time ORDER BY symbol, bar_time;
""")
return cur.fetchall()
ts_agg_time, ts_agg_result = time_read(timescaledb_ohlcv, n_runs=min(3, TIMING_RUNS))
validate_result(ts_agg_result, agg_expected_rows, "TimescaleDB aggregation")
results.append(
BenchmarkResult("TimescaleDB", "aggregation", ts_agg_time, 0, len(ts_agg_result))
)
print(f"\nTimescaleDB: {ts_size / 1e6:.1f} MB")
print(f" Write: {ts_write_time:.3f}s | Read: {ts_read_time:.3f}s")
report_row_build("TimescaleDB", ts_write_time)
print(f" Range query: {ts_range_time:.3f}s ({len(ts_range_result):,} rows)")
print(f" OHLCV aggregation: {ts_agg_time:.3f}s ({len(ts_agg_result):,} daily bars)")
cur.execute("DROP TABLE IF EXISTS ohlcv_benchmark CASCADE;")
cur.close()
conn.close()
else:
print("\n⊘ TimescaleDB benchmark skipped")
# %% [markdown]
# ## PostgreSQL (Vanilla RDBMS Baseline)
#
# Vanilla PostgreSQL serves as the relational baseline: same SQL surface as
# TimescaleDB but without hypertables, compression, or time-series functions.
# The comparison isolates what TimescaleDB's time-series extensions buy you
# on the same engine.
# %%
if HAS_POSTGRES:
print("\n" + "=" * 70)
print("POSTGRESQL BENCHMARK")
print("=" * 70)
benchmark_status["PostgreSQL"]["tested"] = True
from psycopg2.extras import execute_values
pg_conn = psycopg2.connect(
host=DB_CONFIG["postgres"]["host"],
port=DB_CONFIG["postgres"]["port"],
user=DB_CONFIG["postgres"]["user"],
password=DB_CONFIG["postgres"]["password"],
database=DB_CONFIG["postgres"]["database"],
)
pg_conn.autocommit = True
pg_cur = pg_conn.cursor()
pg_cur.execute("DROP TABLE IF EXISTS ohlcv_benchmark CASCADE;")
pg_cur.execute("""
CREATE TABLE ohlcv_benchmark (
timestamp TIMESTAMPTZ NOT NULL, symbol TEXT NOT NULL,
open DOUBLE PRECISION, high DOUBLE PRECISION, low DOUBLE PRECISION,
close DOUBLE PRECISION, volume BIGINT, vwap DOUBLE PRECISION, num_trades INTEGER
);
""")
pg_cur.execute("CREATE INDEX idx_pg_symbol_timestamp ON ohlcv_benchmark(symbol, timestamp);")
# %% [markdown]
# ### PostgreSQL write
#
# The same `execute_values` interface as TimescaleDB, so the same per-row Python
# construction, built the same way and measured separately for the same reason.
# %%
if HAS_POSTGRES:
def build_postgres_rows():
return [
(
r.timestamp,
r.symbol,
r.open,
r.high,
r.low,
r.close,
int(r.volume),
r.vwap,
int(r.num_trades),
)
for r in ohlcv_pandas.itertuples(index=False)
]
def write_postgres():
execute_values(
pg_cur,
"""
INSERT INTO ohlcv_benchmark (timestamp, symbol, open, high, low, close, volume, vwap, num_trades) VALUES %s
""",
build_postgres_rows(),
)
pg_write_time, _ = time_write(write_postgres)
measure_row_build("PostgreSQL", build_postgres_rows)
pg_cur.execute("SELECT pg_total_relation_size('ohlcv_benchmark');")
pg_size = pg_cur.fetchone()[0]
results.append(BenchmarkResult("PostgreSQL", "write", pg_write_time, pg_size, total_rows))
# %%
if HAS_POSTGRES:
# Read
def read_postgres():
pg_cur.execute("SELECT * FROM ohlcv_benchmark;")
return pg_cur.fetchall()
pg_read_time, pg_result = time_read(read_postgres, n_runs=min(3, TIMING_RUNS))
validate_result(pg_result, total_rows, "PostgreSQL read")
results.append(BenchmarkResult("PostgreSQL", "read", pg_read_time, pg_size, total_rows))
# %%
if HAS_POSTGRES:
# Range query
def postgres_range_query():
pg_cur.execute(
"SELECT * FROM ohlcv_benchmark WHERE timestamp >= %s AND timestamp < %s;",
(range_start, range_end),
)
return pg_cur.fetchall()
pg_range_time, pg_range_result = time_read(postgres_range_query, n_runs=min(3, TIMING_RUNS))
validate_result(pg_range_result, range_expected_rows, "PostgreSQL range query")
results.append(
BenchmarkResult("PostgreSQL", "range_query", pg_range_time, 0, len(pg_range_result))
)
# %%
if HAS_POSTGRES:
# Aggregation (date_trunc to the day; emulates time_bucket without the TimescaleDB extension)
def postgres_ohlcv():
pg_cur.execute("""
SELECT symbol, date_trunc('day', timestamp) as bar_time,
(array_agg(open ORDER BY timestamp))[1] as open,
max(high) as high, min(low) as low,
(array_agg(close ORDER BY timestamp DESC))[1] as close,
sum(volume) as volume
FROM ohlcv_benchmark GROUP BY symbol, bar_time ORDER BY symbol, bar_time;
""")
return pg_cur.fetchall()
pg_agg_time, pg_agg_result = time_read(postgres_ohlcv, n_runs=min(3, TIMING_RUNS))
validate_result(pg_agg_result, agg_expected_rows, "PostgreSQL aggregation")
results.append(BenchmarkResult("PostgreSQL", "aggregation", pg_agg_time, 0, len(pg_agg_result)))
print(f"\nPostgreSQL: {pg_size / 1e6:.1f} MB")
print(f" Write: {pg_write_time:.3f}s | Read: {pg_read_time:.3f}s")
report_row_build("PostgreSQL", pg_write_time)
print(f" Range query: {pg_range_time:.3f}s ({len(pg_range_result):,} rows)")
print(f" OHLCV aggregation: {pg_agg_time:.3f}s ({len(pg_agg_result):,} daily bars)")
pg_cur.execute("DROP TABLE IF EXISTS ohlcv_benchmark CASCADE;")
pg_cur.close()
pg_conn.close()
else:
print("\n⊘ PostgreSQL benchmark skipped")
# %% [markdown]
# ## InfluxDB (Time-Series Database)
#
# InfluxDB is a purpose-built time-series database with a tag-based data model
# and the Flux query language. We benchmark line-protocol writes via the Python
# client, full-bucket Flux reads, and per-minute aggregation via
# `aggregateWindow` joined back with `pivot`.
# %%
if HAS_INFLUXDB:
print("\n" + "=" * 70)
print("INFLUXDB BENCHMARK")
print("=" * 70)
benchmark_status["InfluxDB"]["tested"] = True
from influxdb_client import InfluxDBClient, Point, WritePrecision
from influxdb_client.client.write_api import SYNCHRONOUS
influx_url = f"http://{DB_CONFIG['influxdb']['host']}:{DB_CONFIG['influxdb']['port']}"
influx_org = DB_CONFIG["influxdb"]["org"]
influx_token = DB_CONFIG["influxdb"]["token"]
influx_bucket = DB_CONFIG["influxdb"]["bucket"]
# Read timeout in ms: the L-scale (1M-row) Flux read + aggregate/pivot queries
# take well over the client default, so allow several minutes before giving up.
influx_client = InfluxDBClient(
url=influx_url, token=influx_token, org=influx_org, timeout=600_000
)
# Recreate the bucket for a clean run
buckets_api = influx_client.buckets_api()
existing = buckets_api.find_bucket_by_name(influx_bucket)
if existing is not None:
buckets_api.delete_bucket(existing)
orgs = influx_client.organizations_api().find_organizations()
org = next((o for o in orgs if o.name == influx_org), None)
if org is None:
raise RuntimeError(f"InfluxDB org {influx_org!r} not found on the server")
buckets_api.create_bucket(bucket_name=influx_bucket, org_id=org.id)
del buckets_api, existing, orgs, org
# %% [markdown]
# ### InfluxDB write
#
# The line-protocol client takes one `Point` per row, so InfluxDB pays a per-row
# Python cost like PostgreSQL and TimescaleDB, and like theirs it is inside the timed
# write and measured separately. The poll to first-queryable is inside the timed
# region too, so this write ends at the same event as every other one: the rows are
# readable. A fixed sleep outside the timed call would charge InfluxDB nothing for
# the acknowledge-early behaviour that makes it fast.
# %%
if HAS_INFLUXDB and benchmark_status["InfluxDB"]["tested"]:
INFLUX_BATCH_ROWS = 10_000
influx_write_api = influx_client.write_api(write_options=SYNCHRONOUS)
_influx_ts = pd.to_datetime(ohlcv_pandas["timestamp"])
if _influx_ts.dt.tz is None:
_influx_ts = _influx_ts.dt.tz_localize("UTC")
else:
_influx_ts = _influx_ts.dt.tz_convert("UTC")
influx_ts_pandas = _influx_ts
def build_influx_points():
return [
Point("ohlcv")
.tag("symbol", row.symbol)
.field("open", float(row.open))
.field("high", float(row.high))
.field("low", float(row.low))
.field("close", float(row.close))
.field("volume", int(row.volume))
.field("vwap", float(row.vwap))
.field("num_trades", int(row.num_trades))
.time(ts, WritePrecision.NS)
for ts, row in zip(influx_ts_pandas, ohlcv_pandas.itertuples(index=False), strict=True)
]
def write_influxdb():
points = build_influx_points()
# Batched so memory and request size stay bounded.
for i in range(0, len(points), INFLUX_BATCH_ROWS):
influx_write_api.write(
bucket=influx_bucket, org=influx_org, record=points[i : i + INFLUX_BATCH_ROWS]
)
def influx_row_count() -> int:
count_flux = f"""
from(bucket: "{influx_bucket}")
|> range(start: 0)
|> filter(fn: (r) => r._measurement == "ohlcv" and r._field == "close")
|> count()
|> group()
|> sum()
"""
try:
tables = influx_query_api_probe.query(count_flux, org=influx_org)
except Exception:
return 0
return int(sum(r.get_value() for t in tables for r in t.records))
influx_query_api_probe = influx_client.query_api()
def write_influxdb_durable():
write_influxdb()
return wait_until_rows_visible(
influx_row_count, total_rows, timeout=max(120.0, WAL_FLUSH_TIMEOUT * 40)
)
influx_write_time, influx_visible = time_write(write_influxdb_durable)
assert influx_visible == total_rows, f"InfluxDB ingested {influx_visible:,} of {total_rows:,}"
measure_row_build("Influبا ذکر منبع و مطابق مجوز اثر، بهطور کامل نمایش داده میشود. مجوز: MIT
این خلاصه را عامل پژوهشی Stratmill بر پایه متن اصلی نوشته است؛ نسخهای از اثر منبع نیست.