Timing benchmark: where the time goes¶
Every operation ancestree performs, measured. Store lifecycle, node creation, chunking throughput, ingest by file type and size, deduplication hits, reads cold and warm, the query surface at scale, maintenance, and graph rendering.
The point is not a single headline number. It is knowing which operations cost microseconds, which cost milliseconds, and which scale with your data so you can put them in the right place in a pipeline.
Everything runs in a temp directory and is deleted at the end. Absolute numbers depend on the machine; the ratios are what transfer.
import json
import platform
import shutil
import statistics
import sys
import tempfile
import time
import zlib
from pathlib import Path
import matplotlib.pyplot as plt
import numpy as np
import pandas as pd
import ancestree
from ancestree.ingest.cdc import (
AVG_SIZE,
LARGE_FILE_THRESHOLD,
MAX_SIZE,
MIN_SIZE,
chunk_bytes,
delta_encode,
super_features,
)
WORKDIR = Path(tempfile.mkdtemp(prefix="ancestree-timing-"))
MB = 1024 * 1024
print(f"ancestree {ancestree.__version__}")
print(f"python {sys.version.split()[0]}")
print(f"platform {platform.platform()}")
print(f"machine {platform.processor() or platform.machine()}")
print(f"workdir {WORKDIR}")
print(
f"\nchunker min {MIN_SIZE // 1024}K avg {AVG_SIZE // 1024}K "
f"max {MAX_SIZE // 1024}K large-file fallback at "
f"{LARGE_FILE_THRESHOLD // MB}M"
)
ancestree 0.2.0 python 3.12.12 platform macOS-26.5.1-arm64-arm-64bit machine arm workdir /var/folders/xf/n_m7ztrx4x577r3935n_1m1w0000gn/T/ancestree-timing-y3jm8yo4 chunker min 8K avg 16K max 256K large-file fallback at 64M
The harness¶
Every measurement is a median of repeated runs after a warm-up, so a single
scheduling hiccup cannot set the number. bench returns milliseconds.
RESULTS: list[dict] = []
def bench(label, fn, reps=7, warmup=1, group="", note=""):
"""Median wall time of `fn` in ms, recorded in RESULTS."""
for _ in range(warmup):
fn()
samples = []
for _ in range(reps):
started = time.perf_counter()
fn()
samples.append((time.perf_counter() - started) * 1000)
ms = statistics.median(samples)
RESULTS.append(
{
"group": group,
"operation": label,
"ms": ms,
"spread_ms": max(samples) - min(samples),
"note": note,
}
)
print(f" {label:<44} {ms:9.3f} ms (+/- {max(samples) - min(samples):.3f})")
return ms
def once(label, fn, group="", note=""):
"""A single timed run, for operations that cannot repeat (they consume
the state they measure): store creation, prune, the first cold open."""
started = time.perf_counter()
result = fn()
ms = (time.perf_counter() - started) * 1000
RESULTS.append(
{"group": group, "operation": label, "ms": ms, "spread_ms": 0.0, "note": note}
)
print(f" {label:<44} {ms:9.3f} ms")
return result, ms
The corpus¶
Eight payload generators spanning what actually lands in a pipeline: text formats, numeric arrays, and already-compressed blobs the store will decline to compress again. Each returns bytes at roughly the size asked for, so throughput numbers across types are comparable.
rng = np.random.default_rng(0)
def make_csv(target_bytes):
rows = max(1, target_bytes // 22)
return (
pd.DataFrame(
{
"ts": np.arange(rows),
"a": rng.normal(size=rows).round(4),
"b": rng.normal(size=rows).round(4),
}
)
.to_csv(index=False)
.encode()
)
def make_jsonl(target_bytes):
rows = max(1, target_bytes // 78)
return b"".join(
json.dumps(
{
"id": i,
"user": f"user_{i % 500}",
"score": round(float(rng.random()), 6),
"ok": bool(i % 3),
"tag": "abcdefghij"[i % 10] * 8,
}
).encode()
+ b"\n"
for i in range(rows)
)
def make_log(target_bytes):
levels = ["INFO", "WARN", "ERROR", "DEBUG"]
rows = max(1, target_bytes // 86)
return b"".join(
f"2026-07-25T12:{i % 60:02d}:{i % 60:02d}Z {levels[i % 4]:5} worker-{i % 16:02d} "
f"processed batch {i} in {rng.random() * 100:.3f}ms\n".encode()
for i in range(rows)
)
def make_float64(target_bytes):
return rng.normal(size=max(1, target_bytes // 8)).tobytes()
def make_int32(target_bytes):
return rng.integers(
0, 1000, size=max(1, target_bytes // 4), dtype=np.int32
).tobytes()
def make_sparse(target_bytes):
"""Mostly zeros: the case zlib demolishes."""
arr = np.zeros(max(1, target_bytes // 8))
idx = rng.choice(arr.size, size=max(1, arr.size // 100), replace=False)
arr[idx] = rng.normal(size=idx.size)
return arr.tobytes()
def make_random(target_bytes):
"""Incompressible: stands in for PNG, parquet, zip, any packed format."""
return rng.integers(0, 256, size=target_bytes, dtype=np.uint8).tobytes()
def make_repetitive(target_bytes):
"""Pathological for content-defined chunking: no byte variety means the
rolling hash almost never fires, so every chunk runs to MAX_SIZE."""
return (
b"the quick brown fox jumps over the lazy dog. " * (target_bytes // 44 + 1)
)[:target_bytes]
GENERATORS = {
"csv": make_csv,
"jsonl": make_jsonl,
"log": make_log,
"float64": make_float64,
"int32": make_int32,
"sparse": make_sparse,
"random": make_random,
"repetitive": make_repetitive,
}
for name, gen in GENERATORS.items():
sample = gen(1 * MB)
print(
f" {name:<11} {len(sample) / MB:5.2f} MB "
f"zlib-6 ratio {len(sample) / max(len(zlib.compress(sample, 6)), 1):6.2f}x"
)
csv 0.93 MB zlib-6 ratio 2.63x jsonl 1.07 MB zlib-6 ratio 7.22x log 0.81 MB zlib-6 ratio 7.21x float64 1.00 MB zlib-6 ratio 1.04x int32 1.00 MB zlib-6 ratio 2.23x sparse 1.00 MB zlib-6 ratio 68.05x random 1.00 MB zlib-6 ratio 1.00x repetitive 1.00 MB zlib-6 ratio 336.30x
1. Store lifecycle¶
What it costs to open the door. Creating a store runs the schema; reopening an existing one reads the config row and sweeps for scratch directories orphaned by a crashed session. Neither replays an index, which is why a cold open does not care how many nodes are in the store.
print("Store lifecycle")
lifecycle_root = WORKDIR / "lifecycle"
_, create_ms = once(
"LineageStore(...) create empty store",
lambda: ancestree.LineageStore(lifecycle_root),
group="lifecycle",
)
seed = ancestree.LineageStore(lifecycle_root)
for i in range(200):
with seed.create_node(step_type="seed") as node:
node.add_meta("i", i)
seed.close()
def reopen():
store = ancestree.LineageStore(lifecycle_root)
store.close()
bench("open + close existing store (200 nodes)", reopen, reps=5, group="lifecycle")
Store lifecycle LineageStore(...) create empty store 2.829 ms open + close existing store (200 nodes) 0.183 ms (+/- 0.059)
0.1830000001064036
2. Node creation¶
A metadata-only node with no artifacts is the floor of the write path: the
with block, provenance capture, and one transaction. Provenance is timed
separately because it is the single largest component and the only one that
leaves the process (two git subprocesses, run concurrently).
print("Node creation")
node_root = WORKDIR / "nodes"
nstore = ancestree.LineageStore(node_root)
counter = iter(range(1_000_000))
def make_node(meta_entries=1, parent=None):
with nstore.create_node(step_type="step", parent=parent) as node:
for k in range(meta_entries):
node.add_meta(f"key_{k}", next(counter))
return node.node_id
bench("create_node, 1 metadata entry", lambda: make_node(1), reps=15, group="write")
bench("create_node, 10 metadata entries", lambda: make_node(10), reps=15, group="write")
bench("create_node, 50 metadata entries", lambda: make_node(50), reps=15, group="write")
from ancestree.domain import provenance
bench(
" of which: provenance capture (2x git)",
provenance.capture,
reps=15,
group="write",
note="component",
)
Node creation
create_node, 1 metadata entry 8.389 ms (+/- 1.646)
create_node, 10 metadata entries 8.413 ms (+/- 0.814)
create_node, 50 metadata entries 8.831 ms (+/- 0.622)
of which: provenance capture (2x git) 7.257 ms (+/- 1.360)
7.256708999193506
Metadata entries are cheap and linear; the fixed cost of the block dominates
until you attach dozens of them. Provenance is the floor you cannot avoid,
and it is deliberately not cached: the worktree can be committed between two
nodes, and a stale git_dirty would misreport reproducibility.
per_entry = [r["ms"] for r in RESULTS if r["operation"].startswith("create_node, ")]
fixed = per_entry[0] - (per_entry[1] - per_entry[0]) / 9
prov = next(r["ms"] for r in RESULTS if "provenance" in r["operation"])
print(f" fixed cost per node ~{fixed:6.2f} ms")
print(f" of which provenance ~{prov:6.2f} ms ({prov / fixed:.0%})")
print(
f" marginal cost per metadata ~{(per_entry[2] - per_entry[0]) / 49 * 1000:6.1f} us"
)
fixed cost per node ~ 8.39 ms of which provenance ~ 7.26 ms (87%) marginal cost per metadata ~ 9.0 us
3. Chunking throughput¶
Below the store, the Gear rolling hash walks the bytes looking for content-defined boundaries. This is pure Python and it is the reason ingest scales with megabytes rather than with file count. Timed here in isolation, with no store, no compression and no database.
print("Chunking throughput (chunk_bytes, 4 MB payloads)")
chunk_rows = []
for name, gen in GENERATORS.items():
payload = gen(4 * MB)
ms = bench(
f"chunk_bytes: {name}",
lambda p=payload: list(chunk_bytes(p)),
reps=3,
group="chunking",
)
parts = list(chunk_bytes(payload))
chunk_rows.append(
{
"type": name,
"MB/s": len(payload) / MB / (ms / 1000),
"chunks": len(parts),
"avg_chunk_KiB": len(payload) / len(parts) / 1024,
}
)
chunk_df = pd.DataFrame(chunk_rows).sort_values("MB/s", ascending=False)
print()
print(chunk_df.to_string(index=False, float_format=lambda v: f"{v:.2f}"))
Chunking throughput (chunk_bytes, 4 MB payloads)
chunk_bytes: csv 148.108 ms (+/- 0.047)
chunk_bytes: jsonl 218.404 ms (+/- 2.159)
chunk_bytes: log 108.354 ms (+/- 0.177)
chunk_bytes: float64 157.962 ms (+/- 1.494)
chunk_bytes: int32 160.669 ms (+/- 0.435)
chunk_bytes: sparse 285.567 ms (+/- 0.129)
chunk_bytes: random 159.499 ms (+/- 0.633)
chunk_bytes: repetitive 292.108 ms (+/- 0.071)
type MB/s chunks avg_chunk_KiB
log 30.34 204 16.50
csv 26.02 209 18.88
float64 25.32 215 19.05
random 25.08 210 19.50
int32 24.90 208 19.69
jsonl 19.68 160 27.50
sparse 14.01 22 186.18
repetitive 13.69 16 256.00
Real data sits in a tight band around 25 MB/s. The outliers are sparse and
repetitive, and they are slow for the same reason: with too little byte
variety the fingerprint almost never crosses the boundary mask, so every
chunk runs the full 256 KiB to the hard cut.
That is a throughput penalty, not a saving. The Gear loop skips the first
MIN_SIZE bytes of every chunk, because no boundary may fall there. At a
16 KiB average that skip covers half the payload; at a 256 KiB chunk it
covers 3%, so almost every byte gets hashed. Degenerate data is both slower
to chunk and worse to store, because 16 chunks per 4 MB is far too coarse
to share anything. Watch the chunk count, not the megabytes per second.
fig, ax = plt.subplots(figsize=(8, 3.5))
colors = ["#c44" if t == "repetitive" else "#468" for t in chunk_df["type"]]
ax.barh(chunk_df["type"], chunk_df["MB/s"], color=colors)
ax.set_xlabel("MB/s")
ax.set_title("Chunking throughput by payload type (red = pathological)")
ax.invert_yaxis()
fig.tight_layout()
plt.show()
The rest of the ingest pipeline¶
Chunking is one of four stages. Each chunk is also hashed, compressed, and fingerprinted for the resemblance index that Layer 2 uses to find delta bases. Timed per megabyte on a single representative payload.
print("Ingest stages, per 4 MB of CSV")
csv_payload = make_csv(4 * MB)
csv_chunks = list(chunk_bytes(csv_payload))
import hashlib
bench(
"stage: chunk (Gear boundaries)",
lambda: list(chunk_bytes(csv_payload)),
reps=3,
group="stages",
)
bench(
"stage: sha256 every chunk",
lambda: [hashlib.sha256(c).hexdigest() for c in csv_chunks],
reps=3,
group="stages",
)
bench(
"stage: zlib-6 every chunk",
lambda: [zlib.compress(c, 6) for c in csv_chunks],
reps=3,
group="stages",
)
bench(
"stage: super_features every chunk",
lambda: [super_features(c) for c in csv_chunks],
reps=3,
group="stages",
)
bench(
"delta_encode one chunk against a base",
lambda: delta_encode(csv_chunks[0], csv_chunks[1]),
reps=15,
group="stages",
)
Ingest stages, per 4 MB of CSV stage: chunk (Gear boundaries) 148.965 ms (+/- 0.256) stage: sha256 every chunk 1.467 ms (+/- 0.125) stage: zlib-6 every chunk 147.810 ms (+/- 0.255) stage: super_features every chunk 17.150 ms (+/- 0.055) delta_encode one chunk against a base 0.806 ms (+/- 0.128)
0.8062920005613705
4. Write path by file type¶
The full round trip: what your code pays writing the file natively, versus
what create_node costs end to end. The difference is the packing that
happens at block exit, and it is the price of everything the store gives you.
Each type goes into a fresh store, written once. Repeating a write into the same store would measure a deduplication hit rather than a cold ingest, which is a different (and much faster) number, measured separately in section 5.
print("Write: direct vs through the store (4 MB payloads, cold)")
direct_dir = WORKDIR / "write_direct"
direct_dir.mkdir(parents=True)
write_rows = []
for name, gen in GENERATORS.items():
payload = gen(4 * MB)
target = direct_dir / f"{name}.bin"
direct = bench(
f"direct write: {name}",
lambda p=payload, t=target: t.write_bytes(p),
reps=5,
group="write-direct",
)
wstore = ancestree.LineageStore(WORKDIR / f"write_{name}", reuse_identical=False)
def store_write(p=payload, n=name, s=wstore):
with s.create_node(step_type="w") as node:
(node / f"{n}.bin").write_bytes(p)
node.add_meta("i", next(counter))
through = bench(
f"store write: {name}", store_write, reps=1, warmup=0, group="write-store"
)
wstore.close()
write_rows.append(
{
"type": name,
"direct_ms": direct,
"store_ms": through,
"overhead_x": through / direct,
"store_MB/s": 4 / (through / 1000),
}
)
write_df = pd.DataFrame(write_rows).sort_values("store_MB/s", ascending=False)
print()
print(write_df.to_string(index=False, float_format=lambda v: f"{v:.2f}"))
Write: direct vs through the store (4 MB payloads, cold)
direct write: csv 0.741 ms (+/- 0.407)
store write: csv 352.487 ms (+/- 0.000)
direct write: jsonl 0.841 ms (+/- 0.353)
store write: jsonl 321.678 ms (+/- 0.000)
direct write: log 0.602 ms (+/- 0.034)
store write: log 215.256 ms (+/- 0.000)
direct write: float64 0.687 ms (+/- 0.031)
store write: float64 262.726 ms (+/- 0.000)
direct write: int32 0.685 ms (+/- 0.043)
store write: int32 552.748 ms (+/- 0.000)
direct write: sparse 0.676 ms (+/- 0.059)
store write: sparse 334.900 ms (+/- 0.000)
direct write: random 0.693 ms (+/- 0.048)
store write: random 244.386 ms (+/- 0.000)
direct write: repetitive 0.690 ms (+/- 0.031)
store write: repetitive 346.444 ms (+/- 0.000)
type direct_ms store_ms overhead_x store_MB/s
log 0.60 215.26 357.62 18.58
random 0.69 244.39 352.63 16.37
float64 0.69 262.73 382.40 15.23
jsonl 0.84 321.68 382.55 12.43
sparse 0.68 334.90 495.35 11.94
repetitive 0.69 346.44 501.94 11.55
csv 0.74 352.49 475.56 11.35
int32 0.69 552.75 806.54 7.24
The multiplier looks alarming because a direct write of a few megabytes is nearly free: the OS takes the bytes and returns. What matters is the absolute figure. Packing runs at tens of MB/s, so it is felt on artifacts measured in hundreds of megabytes, and invisible on the metadata-heavy nodes that make up most of a pipeline.
Write cost against file size¶
Whether it is linear, measured into a fresh store each time so the growing chunk pool of the previous run cannot contaminate the next. The 64 MiB threshold is included: above it the chunker stops looking for content-defined boundaries and cuts at fixed offsets at C speed instead.
print("Write cost against artifact size (CSV, fresh store each time)")
size_rows = []
for mb in [0.25, 1, 4, 16, 48, 80]:
payload = make_csv(int(mb * MB))
sstore = ancestree.LineageStore(WORKDIR / f"size_{mb}", reuse_identical=False)
def store_write(p=payload, s=sstore):
with s.create_node(step_type="sz") as node:
(node / "data.csv").write_bytes(p)
node.add_meta("i", next(counter))
ms = bench(
f"write {mb:>5} MB artifact", store_write, reps=1, warmup=0, group="write-size"
)
sstore.close()
size_rows.append(
{
"MB": mb,
"ms": ms,
"MB/s": mb / (ms / 1000),
"path": "fixed cuts" if mb * MB >= LARGE_FILE_THRESHOLD else "CDC",
}
)
size_df = pd.DataFrame(size_rows)
print()
print(size_df.to_string(index=False, float_format=lambda v: f"{v:.2f}"))
Write cost against artifact size (CSV, fresh store each time) write 0.25 MB artifact 32.487 ms (+/- 0.000) write 1 MB artifact 91.099 ms (+/- 0.000) write 4 MB artifact 353.378 ms (+/- 0.000) write 16 MB artifact 1519.735 ms (+/- 0.000) write 48 MB artifact 5009.462 ms (+/- 0.000) write 80 MB artifact 6287.517 ms (+/- 0.000) MB ms MB/s path 0.25 32.49 7.70 CDC 1.00 91.10 10.98 CDC 4.00 353.38 11.32 CDC 16.00 1519.74 10.53 CDC 48.00 5009.46 9.58 CDC 80.00 6287.52 12.72 fixed cuts
fig, ax = plt.subplots(figsize=(7, 3.5))
ax.plot(size_df["MB"], size_df["ms"], "o-", color="#468")
ax.axvline(
LARGE_FILE_THRESHOLD / MB,
color="#c44",
ls="--",
label=f"large-file threshold ({LARGE_FILE_THRESHOLD // MB} MiB)",
)
ax.set_xlabel("artifact size (MB)")
ax.set_ylabel("write time (ms)")
ax.set_title("Ingest cost against artifact size")
ax.legend()
fig.tight_layout()
plt.show()
Roughly linear, at a throughput that improves with size as the fixed cost of the node amortises away. Ingest is a megabytes-per-second budget: a quarter-megabyte artifact is free, a 16 MB one is noticeable, and a hundred-megabyte one is something to put at a step boundary rather than inside a loop.
5. What deduplication costs and saves in time¶
Two different hits. A node-level hit is a whole rerun that produced identical content: the store recognises it and stores nothing. A chunk-level hit is a distinct node whose bytes are already in the pool: the node is written, the chunks are not.
Both are compared against a cold write of the same payload into a store that has never seen it.
print("Deduplication, in time")
dedup_payload = make_csv(8 * MB)
cold_store = ancestree.LineageStore(WORKDIR / "dedup_cold", reuse_identical=True)
def cold_write():
with cold_store.create_node(step_type="d") as node:
(node / "data.csv").write_bytes(dedup_payload)
node.add_meta("run", next(counter))
_, cold_ms = once("cold write, 8 MB (nothing in the pool)", cold_write, group="dedup")
def chunk_hit():
with cold_store.create_node(step_type="d") as node:
(node / "data.csv").write_bytes(dedup_payload)
node.add_meta("run", next(counter))
chunk_ms = bench(
"chunk-level hit (same bytes, new node)", chunk_hit, reps=3, group="dedup"
)
dedup_store = ancestree.LineageStore(WORKDIR / "dedup_node", reuse_identical=True)
with dedup_store.create_node(step_type="d") as first:
(first / "data.csv").write_bytes(dedup_payload)
first.add_meta("fixed", 1)
def node_hit():
with dedup_store.create_node(step_type="d") as node:
(node / "data.csv").write_bytes(dedup_payload)
node.add_meta("fixed", 1)
node_ms = bench("node-level hit (identical rerun)", node_hit, reps=3, group="dedup")
print(f"\n chunk-level hit is {cold_ms / chunk_ms:.2f}x faster than a cold write")
print(f" node-level hit is {cold_ms / node_ms:.2f}x faster than a cold write")
print(
f" nodes in the dedup store after {1 + 3 + 1} identical writes: "
f"{dedup_store.stats()['nodes']}"
)
Deduplication, in time cold write, 8 MB (nothing in the pool) 719.770 ms chunk-level hit (same bytes, new node) 328.165 ms (+/- 2.976) node-level hit (identical rerun) 5.303 ms (+/- 0.192) chunk-level hit is 2.19x faster than a cold write node-level hit is 135.73x faster than a cold write nodes in the dedup store after 5 identical writes: 1
A chunk-level hit still pays to read, chunk and hash the bytes: the store cannot know they are duplicates until it has hashed them. What it saves is compression and the write. A node-level hit additionally skips the insert entirely and rebinds onto the node that already exists.
The cost of Layer 2¶
Delta storage against near-identical revisions is the expensive path: every
chunk that misses exact dedup gets fingerprinted, looked up in the
resemblance index, and trial-encoded against a candidate base. Measured as
ingest time for twelve successive revisions with delta=True and
delta=False.
print("Layer 2 (delta storage) ingest cost")
def revisions(policy, versions=8, size=2 * MB):
root = WORKDIR / f"layer2_{policy}"
store = ancestree.LineageStore(root, delta=policy, reuse_identical=False)
payload = bytearray(make_float64(size))
local = np.random.default_rng(7)
started = time.perf_counter()
for v in range(versions):
for pos in local.integers(0, len(payload), size=len(payload) // 100):
payload[int(pos)] = int(local.integers(0, 256))
with store.create_node(step_type="rev") as node:
(node / "data.bin").write_bytes(bytes(payload))
node.add_meta("v", v)
elapsed = (time.perf_counter() - started) * 1000
stats = store.stats()
store.close()
return elapsed, stats
l1_ms, l1_stats = revisions(False)
l2_ms, l2_stats = revisions(True)
for label, ms, stats in [
("layer 1 only", l1_ms, l1_stats),
("layer 1 + 2", l2_ms, l2_stats),
]:
RESULTS.append(
{
"group": "layer2",
"operation": f"8 revisions, {label}",
"ms": ms,
"spread_ms": 0.0,
"note": "",
}
)
print(
f" {label:<14} {ms:8.1f} ms stored {stats['chunk_stored_bytes'] / MB:6.2f} MB"
f" ratio {stats['dedup_ratio']}"
)
print(
f"\n Layer 2 costs {l2_ms / l1_ms:.2f}x the ingest time and stores "
f"{l1_stats['chunk_stored_bytes'] / l2_stats['chunk_stored_bytes']:.2f}x less."
)
Layer 2 (delta storage) ingest cost layer 1 only 1130.0 ms stored 15.44 MB ratio 1.036 layer 1 + 2 1299.7 ms stored 7.19 MB ratio 2.225 Layer 2 costs 1.15x the ingest time and stores 2.15x less.
6. Read path¶
The first read of an artifact in a session reassembles it from chunks into the session cache. Every read after that is a plain file read. Both are timed against reading the identical bytes from a path with no store involved.
print("Read: direct vs cold reassembly vs warm cache")
read_store = ancestree.LineageStore(WORKDIR / "read", reuse_identical=False)
read_ids = {}
read_direct = {}
for mb in [1, 4, 16]:
payload = make_csv(mb * MB)
plain = direct_dir / f"read_{mb}.csv"
plain.write_bytes(payload)
read_direct[mb] = plain
with read_store.create_node(step_type="r") as node:
(node / "data.csv").write_bytes(payload)
node.add_meta("mb", mb)
read_ids[mb] = node.node_id
read_rows = []
for mb, node_id in read_ids.items():
d = bench(
f"direct read {mb:>2} MB",
lambda p=read_direct[mb]: p.read_bytes(),
reps=5,
group="read",
)
def cold(nid=node_id):
read_store._chunks.clear_cache()
return (read_store.get(nid) / "data.csv").read_bytes()
c = bench(f"cold read {mb:>2} MB (reassemble)", cold, reps=3, group="read")
w = bench(
f"warm read {mb:>2} MB (cache hit)",
lambda nid=node_id: (read_store.get(nid) / "data.csv").read_bytes(),
reps=5,
group="read",
)
read_rows.append(
{
"MB": mb,
"direct_ms": d,
"cold_ms": c,
"warm_ms": w,
"cold_x": c / d,
"warm_x": w / d,
}
)
read_df = pd.DataFrame(read_rows)
print()
print(read_df.to_string(index=False, float_format=lambda v: f"{v:.2f}"))
Read: direct vs cold reassembly vs warm cache direct read 1 MB 0.038 ms (+/- 0.004) cold read 1 MB (reassemble) 3.010 ms (+/- 0.208) warm read 1 MB (cache hit) 0.062 ms (+/- 0.007) direct read 4 MB 0.110 ms (+/- 0.059) cold read 4 MB (reassemble) 10.675 ms (+/- 0.076) warm read 4 MB (cache hit) 0.130 ms (+/- 0.021) direct read 16 MB 0.990 ms (+/- 0.092) cold read 16 MB (reassemble) 42.034 ms (+/- 0.932) warm read 16 MB (cache hit) 0.913 ms (+/- 0.083) MB direct_ms cold_ms warm_ms cold_x warm_x 1 0.04 3.01 0.06 79.20 1.63 4 0.11 10.68 0.13 97.01 1.18 16 0.99 42.03 0.91 42.45 0.92
Reassembly is a decompress and a concatenate, so it costs a small multiple of a plain read and only once per session. The warm number is a plain read plus a node lookup, which is where the residual overhead lives.
7. Queries at scale¶
Searches are answered by indexed SQL, not by scanning nodes. Timed against a 3000-node store with realistic metadata, the size at which a linear implementation would already be painful.
print("Building a 3000-node store, measuring search cost as it grows...")
query_root = WORKDIR / "query"
qstore = ancestree.LineageStore(query_root)
CHECKPOINTS = [250, 500, 1000, 2000, 3000]
started = time.perf_counter()
chain_parent = None
node_ids = []
scale_rows = []
for i in range(3000):
parent = chain_parent if i % 10 else None
with qstore.create_node(
step_type=["ingest", "clean", "model"][i % 3], parent=parent
) as node:
node.add_meta("run_id", i)
node.add_meta("accuracy", round(float(rng.random()), 4))
node.add_meta("dataset", f"ds_{i % 25}")
node.add_meta("notes", "a routine run", searchable=False)
node_ids.append(node.node_id)
chain_parent = node.node_id
if i + 1 in CHECKPOINTS:
n = i + 1
full = bench(f"find() over {n:>4} nodes", qstore.find, reps=3, group="scale")
sel = bench(
f"find(run_id) in {n:>4} nodes",
lambda k=n: qstore.find(run_id=k - 1),
reps=15,
group="scale",
)
scale_rows.append({"nodes": n, "find_all_ms": full, "selective_ms": sel})
build_s = time.perf_counter() - started
print(f"\n built in {build_s:.1f} s ({build_s / 3000 * 1000:.2f} ms per node)\n")
deep = node_ids[-1]
print("Queries over 3000 nodes")
bench("find() everything", lambda: qstore.find(), reps=5, group="query")
bench(
"find(step_type=...) ~1000 hits",
lambda: qstore.find(step_type="model"),
reps=5,
group="query",
)
bench(
"find(run_id=...) selective, 1 hit",
lambda: qstore.find(run_id=1500),
reps=15,
group="query",
)
bench(
"find(dataset=...) ~120 hits",
lambda: qstore.find(dataset="ds_7"),
reps=15,
group="query",
)
bench(
"find(accuracy=predicate) full scan",
lambda: qstore.find(accuracy=lambda a: a is not None and a > 0.9),
reps=5,
group="query",
)
bench(
"latest(step_type=...)",
lambda: qstore.latest(step_type="clean"),
reps=15,
group="query",
)
bench("get(node_id)", lambda: qstore.get(deep), reps=25, group="query")
bench(
"lineage(node) 10-deep chain", lambda: qstore.lineage(deep), reps=15, group="query"
)
bench("ancestors(node)", lambda: qstore.ancestors(deep), reps=15, group="query")
bench("children(node)", lambda: qstore.children(node_ids[-2]), reps=25, group="query")
bench("stats()", qstore.stats, reps=15, group="query")
bench(
"sql() GROUP BY over node table",
lambda: qstore.sql("SELECT step_type, count(*) FROM node GROUP BY 1"),
reps=15,
group="query",
)
Building a 3000-node store, measuring search cost as it grows... find() over 250 nodes 2.124 ms (+/- 0.025) find(run_id) in 250 nodes 0.022 ms (+/- 0.012) find() over 500 nodes 4.351 ms (+/- 0.282) find(run_id) in 500 nodes 0.022 ms (+/- 0.011) find() over 1000 nodes 8.852 ms (+/- 0.043) find(run_id) in 1000 nodes 0.022 ms (+/- 0.013) find() over 2000 nodes 17.655 ms (+/- 0.043) find(run_id) in 2000 nodes 0.022 ms (+/- 0.014) find() over 3000 nodes 26.860 ms (+/- 22.563) find(run_id) in 3000 nodes 0.023 ms (+/- 0.016) built in 28.3 s (9.42 ms per node) Queries over 3000 nodes find() everything 27.210 ms (+/- 0.355) find(step_type=...) ~1000 hits 9.105 ms (+/- 0.751) find(run_id=...) selective, 1 hit 0.020 ms (+/- 0.009) find(dataset=...) ~120 hits 1.146 ms (+/- 0.090) find(accuracy=predicate) full scan 11.108 ms (+/- 0.109) latest(step_type=...) 2.407 ms (+/- 0.094) get(node_id) 0.011 ms (+/- 0.003) lineage(node) 10-deep chain 0.103 ms (+/- 0.008) ancestors(node) 5.307 ms (+/- 0.612) children(node) 0.019 ms (+/- 0.008) stats() 0.013 ms (+/- 0.009) sql() GROUP BY over node table 0.201 ms (+/- 0.036)
0.2009170002565952
Does search scale?¶
The measurements taken at each checkpoint while the store was being built. A
flat line for the selective lookup is the index doing its job. find() with
no filter has to materialise every node it matches, so it is expected to rise
with the store.
scale_df = pd.DataFrame(scale_rows)
print()
print(scale_df.to_string(index=False, float_format=lambda v: f"{v:.3f}"))
fig, (a1, a2) = plt.subplots(1, 2, figsize=(10, 3.5))
a1.plot(scale_df["nodes"], scale_df["find_all_ms"], "o-", color="#468")
a1.set_title("find() with no filter")
a2.plot(scale_df["nodes"], scale_df["selective_ms"], "o-", color="#484")
a2.set_title("find(run_id=...) selective")
a2.set_ylim(bottom=0)
for ax in (a1, a2):
ax.set_xlabel("nodes in store")
ax.set_ylabel("ms")
fig.tight_layout()
plt.show()
nodes find_all_ms selective_ms 250 2.124 0.022 500 4.351 0.022 1000 8.852 0.022 2000 17.655 0.022 3000 26.860 0.023
8. Maintenance and export¶
The operations you run occasionally rather than in a loop. All of these touch the whole store, so they scale with it: budget for them accordingly.
print("Maintenance (3000-node store)")
bench(
"export() JSON sidecars",
lambda: qstore.export_metadata(WORKDIR / "export_out"),
reps=1,
warmup=0,
group="maintenance",
)
bench(
"backup() online copy",
lambda: qstore.backup(WORKDIR / "backup.db"),
reps=3,
group="maintenance",
)
bench(
"prune(node) dry run",
lambda: qstore.prune(node_ids[1000]),
reps=5,
group="maintenance",
)
once("compact() on a clean store", qstore.compact, group="maintenance")
once(
"prune(deep branch, dry_run=False)",
lambda: qstore.prune(node_ids[1000], dry_run=False),
group="maintenance",
)
once(
"integrity_check via sql()",
lambda: qstore.sql("PRAGMA integrity_check"),
group="maintenance",
)
Maintenance (3000-node store) export() JSON sidecars 515.344 ms (+/- 0.000) backup() online copy 5.872 ms (+/- 0.228) prune(node) dry run 0.225 ms (+/- 0.013) compact() on a clean store 1.981 ms prune(deep branch, dry_run=False) 1.683 ms integrity_check via sql() 13.073 ms
([<sqlite3.Row at 0x11a12a7a0>], 13.073374999294174)
prune defaults to a dry run, which is a graph traversal and nothing else.
The real deletion additionally compacts, which scans the whole chunk pool.
Pruning in a loop should pass compact=False and call compact() once at
the end.
9. Rendering¶
The static snapshot inlines the entire store into one HTML file, so its cost scales with node count and artifact size.
print("Web graph rendering")
for n in [100, 500, 2000]:
render_root = WORKDIR / f"render_{n}"
rstore = ancestree.LineageStore(render_root)
for i in range(n):
with rstore.create_node(step_type="r") as node:
node.add_meta("i", i)
node.add_meta("accuracy", round(float(rng.random()), 4))
path, ms = once(
f"export_graph(), {n:>4} nodes",
lambda s=rstore: s.export_graph(),
group="render",
)
print(
f" -> {path.stat().st_size / 1024:.0f} KiB of HTML "
f"({path.stat().st_size / n:.0f} bytes per node)"
)
rstore.close()
Web graph rendering
export_graph(), 100 nodes 5.614 ms
-> 542 KiB of HTML (5555 bytes per node)
export_graph(), 500 nodes 20.757 ms
-> 839 KiB of HTML (1719 bytes per node)
export_graph(), 2000 nodes 75.520 ms
-> 1953 KiB of HTML (1000 bytes per node)
10. Summary¶
Every measurement in one table, grouped by what it belongs to.
summary = pd.DataFrame(RESULTS)
pd.set_option("display.max_rows", None, "display.width", 140)
for group in summary["group"].unique():
block = summary[summary["group"] == group]
print(f"\n== {group} " + "=" * (70 - len(group)))
print(
block[["operation", "ms"]].to_string(
index=False, float_format=lambda v: f"{v:.3f}"
)
)
== lifecycle =============================================================
operation ms
LineageStore(...) create empty store 2.829
open + close existing store (200 nodes) 0.183
== write =================================================================
operation ms
create_node, 1 metadata entry 8.389
create_node, 10 metadata entries 8.413
create_node, 50 metadata entries 8.831
of which: provenance capture (2x git) 7.257
== chunking ==============================================================
operation ms
chunk_bytes: csv 148.108
chunk_bytes: jsonl 218.404
chunk_bytes: log 108.354
chunk_bytes: float64 157.962
chunk_bytes: int32 160.669
chunk_bytes: sparse 285.567
chunk_bytes: random 159.499
chunk_bytes: repetitive 292.108
== stages ================================================================
operation ms
stage: chunk (Gear boundaries) 148.965
stage: sha256 every chunk 1.467
stage: zlib-6 every chunk 147.810
stage: super_features every chunk 17.150
delta_encode one chunk against a base 0.806
== write-direct ==========================================================
operation ms
direct write: csv 0.741
direct write: jsonl 0.841
direct write: log 0.602
direct write: float64 0.687
direct write: int32 0.685
direct write: sparse 0.676
direct write: random 0.693
direct write: repetitive 0.690
== write-store ===========================================================
operation ms
store write: csv 352.487
store write: jsonl 321.678
store write: log 215.256
store write: float64 262.726
store write: int32 552.748
store write: sparse 334.900
store write: random 244.386
store write: repetitive 346.444
== write-size ============================================================
operation ms
write 0.25 MB artifact 32.487
write 1 MB artifact 91.099
write 4 MB artifact 353.378
write 16 MB artifact 1519.735
write 48 MB artifact 5009.462
write 80 MB artifact 6287.517
== dedup =================================================================
operation ms
cold write, 8 MB (nothing in the pool) 719.770
chunk-level hit (same bytes, new node) 328.165
node-level hit (identical rerun) 5.303
== layer2 ================================================================
operation ms
8 revisions, layer 1 only 1129.953
8 revisions, layer 1 + 2 1299.721
== read ==================================================================
operation ms
direct read 1 MB 0.038
cold read 1 MB (reassemble) 3.010
warm read 1 MB (cache hit) 0.062
direct read 4 MB 0.110
cold read 4 MB (reassemble) 10.675
warm read 4 MB (cache hit) 0.130
direct read 16 MB 0.990
cold read 16 MB (reassemble) 42.034
warm read 16 MB (cache hit) 0.913
== scale =================================================================
operation ms
find() over 250 nodes 2.124
find(run_id) in 250 nodes 0.022
find() over 500 nodes 4.351
find(run_id) in 500 nodes 0.022
find() over 1000 nodes 8.852
find(run_id) in 1000 nodes 0.022
find() over 2000 nodes 17.655
find(run_id) in 2000 nodes 0.022
find() over 3000 nodes 26.860
find(run_id) in 3000 nodes 0.023
== query =================================================================
operation ms
find() everything 27.210
find(step_type=...) ~1000 hits 9.105
find(run_id=...) selective, 1 hit 0.020
find(dataset=...) ~120 hits 1.146
find(accuracy=predicate) full scan 11.108
latest(step_type=...) 2.407
get(node_id) 0.011
lineage(node) 10-deep chain 0.103
ancestors(node) 5.307
children(node) 0.019
stats() 0.013
sql() GROUP BY over node table 0.201
== maintenance ===========================================================
operation ms
export() JSON sidecars 515.344
backup() online copy 5.872
prune(node) dry run 0.225
compact() on a clean store 1.981
prune(deep branch, dry_run=False) 1.683
integrity_check via sql() 13.073
== render ================================================================
operation ms
export_graph(), 100 nodes 5.614
export_graph(), 500 nodes 20.757
export_graph(), 2000 nodes 75.520
What to take away¶
Three cost classes, and they are three orders of magnitude apart.
def pick(op):
return next(r["ms"] for r in RESULTS if r["operation"].startswith(op))
print("Microseconds (free, put them anywhere)")
print(f" get(node_id) {pick('get(node_id)'):8.3f} ms")
print(
f" selective find() {pick('find(run_id=...) selective'):8.3f} ms"
)
print(f" lineage() {pick('lineage(node)'):8.3f} ms")
print("\nMilliseconds (fine per pipeline step, not per row)")
print(f" create_node, metadata only {pick('create_node, 1'):8.3f} ms")
print(f" find() over 3000 nodes {pick('find() everything'):8.3f} ms")
print("\nScales with your data (budget for it)")
print(f" ingest 4 MB of CSV {pick('store write: csv'):8.1f} ms")
print(f" ingest 48 MB artifact {pick('write 48 MB'):8.1f} ms")
print(f" export() 3000 nodes {pick('export()'):8.1f} ms")
Microseconds (free, put them anywhere) get(node_id) 0.011 ms selective find() 0.020 ms lineage() 0.103 ms Milliseconds (fine per pipeline step, not per row) create_node, metadata only 8.389 ms find() over 3000 nodes 27.210 ms Scales with your data (budget for it) ingest 4 MB of CSV 352.5 ms ingest 48 MB artifact 5009.5 ms export() 3000 nodes 515.3 ms
The write path is the only thing that scales with data volume, and it is paid at block exit rather than while your code runs. Everything on the read and query side is indexed and effectively free at the scale a person explores at.
Clean up¶
for store in (nstore, wstore, cold_store, dedup_store, read_store, qstore):
store.close()
shutil.rmtree(WORKDIR, ignore_errors=True)
print("done")
done