Skip to content

Commit 69313bd

Browse files
authored
Merge pull request #2 from SoheilGtex/benchmark/phase3e-experiment
Finalize Phase 3E benchmark study
2 parents c4313f6 + bcfe625 commit 69313bd

8 files changed

Lines changed: 1049 additions & 11 deletions

File tree

‎.github/workflows/ci.yml‎

Lines changed: 10 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -38,3 +38,13 @@ jobs:
3838
run: |
3939
fdp run-all --source synthetic --seed 20270916 --rows 1000 --output-dir "$RUNNER_TEMP/fdp-smoke"
4040
python -c "import json, os; p=os.path.join(os.environ['RUNNER_TEMP'], 'fdp-smoke', 'correctness.json'); d=json.load(open(p)); assert d['checksum_equal'] is True; assert d['database_final_row_count'] == 1000; assert d['integrity_check'] == 'ok'"
41+
- name: Stateful three-strategy correctness smoke
42+
run: |
43+
python scripts/benchmark.py --output-dir "$RUNNER_TEMP/fdp-benchmark" --scales 1000 --repetitions 1
44+
test "$(wc -l < "$RUNNER_TEMP/fdp-benchmark/correctness.csv")" -gt 1
45+
test "$(wc -l < "$RUNNER_TEMP/fdp-benchmark/query_samples.jsonl")" -eq 2250
46+
- name: Regenerate analysis from raw files
47+
run: |
48+
python scripts/analyze_benchmark.py "$RUNNER_TEMP/fdp-benchmark" "$RUNNER_TEMP/fdp-analysis"
49+
test -s "$RUNNER_TEMP/fdp-analysis/summary.csv"
50+
test -s "$RUNNER_TEMP/fdp-analysis/query_summary.csv"

‎Makefile‎

Lines changed: 4 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1,4 +1,4 @@
1-
.PHONY: install dev lint format-check test smoke ci
1+
.PHONY: install dev lint format-check test smoke benchmark-small ci
22

33
install:
44
python -m pip install .
@@ -18,4 +18,7 @@ test:
1818
smoke:
1919
fdp run-all --source synthetic --seed 20270916 --rows 1000 --output-dir .repro/smoke
2020

21+
benchmark-small:
22+
python scripts/benchmark.py --output-dir .repro/benchmark-small --scales 1000 --repetitions 1
23+
2124
ci: lint format-check test

‎README.md‎

Lines changed: 33 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -2,11 +2,12 @@
22

33
Finance Data Pipelines is a small, correctness-first time-series ETL and SQLite loading
44
project. The current release provides a deterministic offline generator, strict validation,
5-
an independent logical-state oracle, and **Strategy A: atomic full replacement**.
5+
an independent logical-state oracle, and three comparable SQLite strategies: **A: atomic
6+
full replacement**, **B: incremental upsert**, and **C: append-only history with a current
7+
projection**.
68

7-
The repository is a baseline for a later controlled systems study. It does not yet contain
8-
Strategies B/C or performance results, and it is not presented as production-ready,
9-
high-performance, scalable, or research-grade software.
9+
The repository includes a reproducible benchmark harness. Its results are cohort-specific and
10+
are not a claim of production readiness, universal performance, scalability, or research novelty.
1011

1112
## Requirements
1213

@@ -58,7 +59,26 @@ verified. It prints correctness fields as JSON and writes:
5859
warehouse.db
5960
```
6061

61-
No throughput, latency, or comparative benchmark conclusion is produced.
62+
The quickstart produces correctness evidence only. Comparative results are generated separately
63+
by the controlled experiment harness described below.
64+
65+
## Controlled experiment
66+
67+
The benchmark harness is `scripts/benchmark.py`. It freezes SQLite WAL and
68+
`synchronous=FULL`, runs stateful workload sequences from one fresh database per strategy and
69+
repetition, rotates strategy order, validates every transition against the independent oracle,
70+
and writes `environment.json`, `dataset_manifest.json`, `runs.jsonl`, `query_samples.jsonl`,
71+
`summary.csv`, `correctness.csv`, and `recovery.csv`. Each accepted condition uses five fixed
72+
queries with two warm-ups and 30 timed warm-cache samples. Analysis is regenerated from raw files
73+
with `scripts/analyze_benchmark.py`. A small offline smoke is reproducible with:
74+
75+
```bash
76+
make benchmark-small
77+
```
78+
79+
The full experiment is intentionally separate from ordinary CI. It uses synthetic data only;
80+
no financial or market conclusion is supported. Any report must identify the exact machine,
81+
source hash, workload definitions, and limitations of its cohort.
6282

6383
## Optional configuration file
6484

@@ -80,7 +100,8 @@ Prices and optional volume use fixed-point integers with scale `10^-6`.
80100
Duplicate/update rules:
81101

82102
- an exact duplicate event is a safe no-op;
83-
- different payloads for the same key and revision reject the entire snapshot;
103+
- different payloads for the same key and revision reject the entire batch, including across
104+
previously committed batches;
84105
- the highest revision is current;
85106
- a lower revision cannot regress the current state;
86107
- an empty replacement is rejected by default.
@@ -120,13 +141,13 @@ only when build and run commands have actually executed in that release environm
120141

121142
## Current maturity and limitations
122143

123-
- Strategy A only; Strategies B/C belong to the next phase.
144+
- Three strategies with shared logical semantics; physical storage trade-offs differ.
124145
- Single-process, single-writer SQLite baseline.
125146
- Synthetic correctness fixture only; no market-behavior claims.
126147
- No concurrency or distributed-system evaluation.
127-
- No performance measurements or rankings.
128-
- Deterministic exception injection tests transaction rollback; they are not a claim that
129-
every OS/power-loss mode has been tested.
148+
- Benchmark results are machine-specific and do not establish universal rankings.
149+
- Deterministic exception injection and process-interruption tests cover transactional recovery;
150+
they are not a claim that every OS/power-loss mode has been tested.
130151

131152
## Project layout
132153

@@ -142,6 +163,8 @@ src/fdp/
142163
encoding.py canonical binary checksum encoding
143164
manifest.py deterministic JSON manifests
144165
parquet_io.py frozen Parquet schema
166+
strategies.py alternative SQLite implementations with shared logical semantics
167+
scripts/ benchmark.py and analyze_benchmark.py
145168
tests/ isolated correctness and CLI tests
146169
.github/workflows/ci.yml
147170
```

‎scripts/analyze_benchmark.py‎

Lines changed: 208 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,208 @@
1+
#!/usr/bin/env python3
2+
"""Regenerate Phase 3E analysis artifacts from raw JSONL/CSV files only."""
3+
4+
from __future__ import annotations
5+
6+
import argparse
7+
import csv
8+
import json
9+
import math
10+
import statistics
11+
from collections import defaultdict
12+
from pathlib import Path
13+
14+
15+
def chart(
16+
path: Path,
17+
title: str,
18+
labels: list[str],
19+
values: dict[tuple[str, str], float],
20+
series: list[str],
21+
) -> None:
22+
width, height, left, top, plot = 1200, 650, 110, 85, 930
23+
maximum = max(values.values() or [1]) * 1.15
24+
colors = ["#4c78a8", "#f58518", "#54a24b"]
25+
svg = [
26+
f'<svg xmlns="http://www.w3.org/2000/svg" width="{width}" height="{height}">',
27+
'<rect width="100%" height="100%" fill="white"/>',
28+
f'<text x="{width / 2}" y="40" text-anchor="middle" font-family="sans-serif" font-size="24" font-weight="bold">{title}</text>', # noqa: E501
29+
]
30+
for i in range(6):
31+
y = top + plot * (1 - i / 5)
32+
value = maximum * i / 5
33+
svg += [
34+
f'<line x1="{left}" y1="{y}" x2="{left + plot}" y2="{y}" stroke="#ddd"/>',
35+
f'<text x="{left - 8}" y="{y + 4}" text-anchor="end" font-family="sans-serif" font-size="12">{value:.0f}</text>', # noqa: E501
36+
]
37+
group = plot / max(1, len(labels))
38+
bar = group / (len(series) + 1)
39+
for li, label in enumerate(labels):
40+
gx = left + li * group
41+
svg.append(
42+
f'<text x="{gx + group / 2}" y="{top + plot + 30}" text-anchor="middle" font-family="sans-serif" font-size="12">{label}</text>' # noqa: E501
43+
)
44+
for si, strategy in enumerate(series):
45+
value = values.get((label, strategy), 0)
46+
bh = value / maximum * plot
47+
svg.append(
48+
f'<rect x="{gx + (si + 0.5) * bar}" y="{top + plot - bh}" width="{bar * 0.78}" height="{bh}" fill="{colors[si % len(colors)]}"><title>{label} {strategy}: {value}</title></rect>' # noqa: E501
49+
)
50+
for si, strategy in enumerate(series):
51+
x = left + plot - 360 + si * 120
52+
svg += [
53+
f'<rect x="{x}" y="{height - 55}" width="14" height="14" fill="{colors[si % len(colors)]}"/>', # noqa: E501
54+
f'<text x="{x + 20}" y="{height - 43}" font-family="sans-serif" font-size="11">{strategy[:14]}</text>', # noqa: E501
55+
]
56+
svg.append("</svg>")
57+
path.write_text("\n".join(svg))
58+
59+
60+
def main() -> int:
61+
parser = argparse.ArgumentParser()
62+
parser.add_argument("raw_dir", type=Path)
63+
parser.add_argument("output_dir", type=Path)
64+
args = parser.parse_args()
65+
args.output_dir.mkdir(parents=True, exist_ok=True)
66+
runs = [
67+
json.loads(line) for line in (args.raw_dir / "runs.jsonl").read_text().splitlines() if line
68+
]
69+
queries = [
70+
json.loads(line)
71+
for line in (args.raw_dir / "query_samples.jsonl").read_text().splitlines()
72+
if line
73+
]
74+
groups = defaultdict(list)
75+
performance_groups = defaultdict(list)
76+
for row in runs:
77+
key = (row["strategy"], row["workload"], row["dataset_size"])
78+
groups[key].append(row)
79+
if row.get("repetition") != 0:
80+
performance_groups[key].append(row)
81+
with (args.output_dir / "summary.csv").open("w", newline="") as handle:
82+
fields = [
83+
"strategy",
84+
"workload",
85+
"dataset_size",
86+
"n",
87+
"median_ms",
88+
"min_ms",
89+
"max_ms",
90+
"iqr_ms",
91+
"median_rows_sec",
92+
"all_correct",
93+
]
94+
writer = csv.DictWriter(handle, fieldnames=fields)
95+
writer.writeheader()
96+
for (strategy, workload, size), rows in sorted(performance_groups.items()):
97+
durations = sorted(float(row["duration_ms"]) for row in rows)
98+
q1 = durations[(len(durations) - 1) // 4]
99+
q3 = durations[3 * (len(durations) - 1) // 4]
100+
writer.writerow(
101+
{
102+
"strategy": strategy,
103+
"workload": workload,
104+
"dataset_size": size,
105+
"n": len(rows),
106+
"median_ms": statistics.median(durations),
107+
"min_ms": min(durations),
108+
"max_ms": max(durations),
109+
"iqr_ms": q3 - q1,
110+
"median_rows_sec": statistics.median(
111+
row["throughput_rows_sec"] for row in rows
112+
),
113+
"all_correct": all(
114+
row["correctness"]["checksum_equal"]
115+
and row["correctness"]["integrity_check"] == "ok"
116+
for row in rows
117+
),
118+
}
119+
)
120+
with (args.output_dir / "correctness.csv").open("w", newline="") as handle:
121+
fields = [
122+
"strategy",
123+
"workload",
124+
"dataset_size",
125+
"runs",
126+
"all_checksums_equal",
127+
"all_integrity_ok",
128+
]
129+
writer = csv.DictWriter(handle, fieldnames=fields)
130+
writer.writeheader()
131+
for key, rows in sorted(groups.items()):
132+
writer.writerow(
133+
{
134+
"strategy": key[0],
135+
"workload": key[1],
136+
"dataset_size": key[2],
137+
"runs": len(rows),
138+
"all_checksums_equal": all(
139+
row["correctness"]["checksum_equal"] for row in rows
140+
),
141+
"all_integrity_ok": all(
142+
row["correctness"]["integrity_check"] == "ok" for row in rows
143+
),
144+
}
145+
)
146+
qgroups = defaultdict(list)
147+
for row in queries:
148+
qgroups[(row["strategy"], row["workload"], row["dataset_size"], row["query"])].append(
149+
float(row["latency_ms"])
150+
)
151+
with (args.output_dir / "query_summary.csv").open("w", newline="") as handle:
152+
fields = ["strategy", "workload", "dataset_size", "query", "n", "median_ms", "p95_ms"]
153+
writer = csv.DictWriter(handle, fieldnames=fields)
154+
writer.writeheader()
155+
for key, values in sorted(qgroups.items()):
156+
values.sort()
157+
writer.writerow(
158+
{
159+
"strategy": key[0],
160+
"workload": key[1],
161+
"dataset_size": key[2],
162+
"query": key[3],
163+
"n": len(values),
164+
"median_ms": statistics.median(values),
165+
"p95_ms": values[max(0, math.ceil(0.95 * len(values)) - 1)],
166+
}
167+
)
168+
labels = sorted({workload for _, workload, _ in performance_groups})
169+
strategies = sorted({strategy for strategy, _, _ in performance_groups})
170+
ingestion = {
171+
(label, strategy): statistics.median(float(row["duration_ms"]) for row in rows)
172+
for (strategy, label, _), rows in performance_groups.items()
173+
}
174+
storage = {
175+
(label, strategy): statistics.median(float(row.get("storage_bytes", 0)) for row in rows)
176+
for (strategy, label, _), rows in performance_groups.items()
177+
}
178+
chart(
179+
args.output_dir / "ingestion_wall_time.svg",
180+
"Median ingestion wall time (ms)",
181+
labels,
182+
ingestion,
183+
strategies,
184+
)
185+
chart(
186+
args.output_dir / "storage_footprint.svg",
187+
"Median storage footprint (bytes)",
188+
labels,
189+
storage,
190+
strategies,
191+
)
192+
query_labels = sorted({row["query"] for row in queries})
193+
qvalues = {
194+
(label, strategy): statistics.median(values)
195+
for (strategy, _, _, label), values in qgroups.items()
196+
}
197+
chart(
198+
args.output_dir / "query_latency.svg",
199+
"Median warm-cache query latency (ms)",
200+
query_labels,
201+
qvalues,
202+
strategies,
203+
)
204+
return 0
205+
206+
207+
if __name__ == "__main__":
208+
raise SystemExit(main())

0 commit comments

Comments
 (0)