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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
887 changes: 887 additions & 0 deletions plans/ADR-007-tallyman-owned-materialization.md

Large diffs are not rendered by default.

721 changes: 721 additions & 0 deletions plans/ADR-008-row-order-of-reads.md

Large diffs are not rendered by default.

453 changes: 453 additions & 0 deletions plans/ADR-009-digest-stability.md

Large diffs are not rendered by default.

117 changes: 117 additions & 0 deletions scripts/spike_bare_read_chaining.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,117 @@
"""ADR-007 evidence: chaining through a bare read of the parent's snapshot, with no xorq cache node anywhere.

Raw xorq plus tallyman's ``classify_build``; no tallyman build machinery. A worthy parent (an aggregate, canonically
sorted) is built with no cache node and materialized by a plain parquet write to ``<snapshots>/<parent_hash>.parquet``.
A child reads that path with ``deferred_read_parquet`` and adds a filter and a computed column.

Questions, in the order printed:

1. What does the child's build hash depend on: the snapshot's path string, or its bytes/mtime?
2. Can a child be built while the parent's snapshot is absent?
3. Does ``load_expr`` of the child's build need the snapshot? What does executing without it raise?
4. After the parent is re-materialized from its own build, does the child execute to the right answer?
5. Is the child classified cheap? (Under today's inlined chaining the same child is worthy.)
6. Does the child's ``expr.yaml`` carry the literal path, so the portable-path placeholder applies to it?
7. Did anything get written under xorq's global cache directory?

uv run python scripts/spike_bare_read_chaining.py
"""

from __future__ import annotations

import hashlib
import os
import tempfile
from pathlib import Path

HOME = Path(tempfile.mkdtemp(prefix="spike_bare_read_"))
os.environ["XORQ_CACHE_DIR"] = str(HOME / "_global_xorq") # must be set before xorq is imported

import numpy as np # noqa: E402
import pyarrow as pa # noqa: E402
import pyarrow.parquet as pq # noqa: E402
import xorq.api as xo # noqa: E402
from xorq.ibis_yaml.compiler import build_expr, load_expr # noqa: E402

from tallyman_xorq.result_cache import classify_build # noqa: E402

N = 400_000
BUILDS = HOME / "builds"
SNAPSHOTS = HOME / "compute_cache" / "result_cache"


def materialize(expr, dest: Path) -> str:
tmp = dest.with_name(f".{dest.name}.{os.getpid()}.tmp")
pq.write_table(expr.to_pyarrow(), tmp, compression="zstd")
os.replace(tmp, dest)
return hashlib.sha256(dest.read_bytes()).hexdigest()[:16]


def child_of(snapshot: Path, schema=None):
parent = xo.deferred_read_parquet(str(snapshot), schema=schema)
return parent.filter(parent.n > 10).mutate(r=parent.s / parent.n)


def attempt(fn) -> str:
try:
return str(fn())
except Exception as exc: # noqa: BLE001 - the spike reports whatever is raised
return f"raises {type(exc).__name__}: {str(exc)[:90]}"


def main() -> None:
SNAPSHOTS.mkdir(parents=True)
source = HOME / "cas" / "deadbeef.parquet"
source.parent.mkdir()
rng = np.random.default_rng(7)
pq.write_table(
pa.table({"id": np.arange(N), "g": rng.integers(0, 20_000, N), "v": rng.integers(0, 1000, N)}), source
)

t = xo.deferred_read_parquet(str(source))
parent = t.group_by("g").agg(n=t.count(), s=t.v.sum()).order_by(["g", "n", "s"])
parent_build = Path(build_expr(parent, builds_dir=BUILDS))
snapshot = SNAPSHOTS / f"{parent_build.name}.parquet"
first_digest = materialize(load_expr(parent_build), snapshot)
print(f"parent {parent_build.name}: {classify_build(parent_build)}, snapshot digest {first_digest}\n")

same_path = Path(build_expr(child_of(snapshot), builds_dir=BUILDS)).name
original = snapshot.read_bytes()
pq.write_table(pa.table({"g": [1, 2], "n": [99, 98], "s": [5, 6]}), snapshot) # other rows, other size and mtime
other_bytes = Path(build_expr(child_of(snapshot), builds_dir=BUILDS)).name
elsewhere = SNAPSHOTS / "0123456789ab.parquet"
elsewhere.write_bytes(original)
other_path = Path(build_expr(child_of(elsewhere), builds_dir=BUILDS)).name
print(f"1. child hash: {same_path}")
print(f" same path with other bytes: {other_bytes}; same bytes at another path: {other_path}")

snapshot.unlink()
print(f"2. compose with snapshot absent, no schema: {attempt(lambda: child_of(snapshot).schema().names)}")
with_schema = attempt(lambda: build_expr(child_of(snapshot, parent.schema()), builds_dir=BUILDS))
print(f" build with snapshot absent, schema given: {with_schema}")

child_build = BUILDS / same_path
print(f"3. load_expr(child) with snapshot absent: {attempt(lambda: type(load_expr(child_build)).__name__)}")
print(f" execute with snapshot absent: {attempt(lambda: load_expr(child_build).count().execute())}")

again = materialize(load_expr(parent_build), snapshot)
rows = int(load_expr(child_build).count().execute())
want = int(parent.filter(parent.n > 10).count().execute())
print(f"4. parent re-materialized from its build: digest unchanged={again == first_digest}")
print(f" child rows {rows}, expected {want}")

inlined = Path(build_expr(parent.filter(parent.n > 10).mutate(r=parent.s / parent.n), builds_dir=BUILDS))
print(f"5. classify_build(bare-read child) = {classify_build(child_build)}")
print(f" classify_build(same child, parent graph inlined) = {classify_build(inlined)}")

yaml_text = (child_build / "expr.yaml").read_text()
inlined_size = len((inlined / "expr.yaml").read_text())
print(f"6. literal snapshot path in child expr.yaml: {str(snapshot) in yaml_text}")
print(f" child expr.yaml is {len(yaml_text):,} bytes; with the parent graph inlined it is {inlined_size:,}")

leaked = sorted(str(f) for f in (HOME / "_global_xorq").rglob("*.parquet"))
print(f"7. parquet files under XORQ_CACHE_DIR: {leaked or 'none'}")


if __name__ == "__main__":
main()
108 changes: 108 additions & 0 deletions scripts/spike_cheap_classifier.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,108 @@
"""ADR-008 evidence (D4): what a test for "cheap" has to look at, now that two kinds of entry stay.

A cheap entry has no file of its own and pages by its parent's ``__row_order``, so the test has to guarantee that the
column stays unique through the plan. D4's first wording was an allow-list of RELATION operations. This script runs
three tests over the same recipe shapes:

- today's deny-list (``result_cache.classify_build``, a regex over ``expr.yaml``);
- a relation-only allow-list, which is D4 as first worded, with drop-null and fill-null added to its list;
- the test D4 now specifies: the relation allow-list, exactly one file read, and no value operation that multiplies
rows, depends on row order, or is not pure, decided by base class on the live expression.

For each shape it also executes the plan and reports whether ``__row_order`` is still unique.

uv run python scripts/spike_cheap_classifier.py
"""

from __future__ import annotations

import os
import re
import tempfile
from pathlib import Path

HOME = Path(tempfile.mkdtemp(prefix="spike_cheap_classifier_"))
os.environ["XORQ_CACHE_DIR"] = str(HOME / "_global_xorq") # must be set before xorq is imported

import pyarrow as pa # noqa: E402
import pyarrow.parquet as pq # noqa: E402
import xorq.api as xo # noqa: E402
import xorq.vendor.ibis as ibis # noqa: E402
import xorq.vendor.ibis.expr.operations as ops # noqa: E402
from xorq.common.utils.graph_utils import walk_nodes # noqa: E402
from xorq.expr.relations import Read # noqa: E402
from xorq.ibis_yaml.compiler import build_expr # noqa: E402
from xorq.vendor.ibis.expr.operations.core import Node # noqa: E402

from tallyman_xorq.result_cache import classify_build # noqa: E402

ROW_PRESERVING_RELATIONS = (Read, ops.Filter, ops.Project, ops.DropColumns, ops.DropNull, ops.FillNull)
NEVER_CHEAP_VALUES = (ops.Unnest, ops.WindowFunction, ops.Impure)
NEVER_CHEAP_BY_NAME = {"TimestampNow", "DateNow"} # these are Constant, not Impure, in xorq's ibis


def relation_only_allow_list(expr) -> bool:
relations = [n for n in walk_nodes((Node,), expr) if isinstance(n, ops.Relation)]
return all(isinstance(n, ROW_PRESERVING_RELATIONS) for n in relations)


def is_cheap(expr) -> bool:
nodes = list(walk_nodes((Node,), expr))
relations = [n for n in nodes if isinstance(n, ops.Relation)]
if len({n for n in relations if isinstance(n, Read)}) != 1:
return False
if not all(isinstance(n, ROW_PRESERVING_RELATIONS) for n in relations):
return False
for n in nodes:
if isinstance(n, NEVER_CHEAP_VALUES) or type(n).__name__ in NEVER_CHEAP_BY_NAME:
return False
if any("UDF" in base.__name__ for base in type(n).__mro__):
return False
return True


def main() -> None:
n = 6
tags = [["x", "y"], ["z"], [], ["x"], ["y", "z", "w"], ["q"]]
parent = {"k": list(range(n)), "v": [1.5, 2.5, None, 4.5, 5.5, 6.5], "tags": tags, "__row_order": list(range(n))}
pq.write_table(pa.table(parent), HOME / "t.parquet")
pq.write_table(pa.table({"k": [1, 3, 5], "__row_order": [0, 1, 2]}), HOME / "u.parquet")
t = xo.deferred_read_parquet(str(HOME / "t.parquet"))
u = xo.deferred_read_parquet(str(HOME / "u.parquet"))

shapes = {
"filter + computed column": t.filter(t.k > 0).mutate(w=t.v * 2),
"rename, cast, drop a column": t.rename(key="k").mutate(v=t.v.cast("float32")).drop("tags"),
"drop_null, fill_null": t.drop_null(["v"]).fill_null({"v": 0.0}),
"unnest inside a select": t.select("k", "__row_order", tag=t.tags.unnest()),
"row_number() in a mutate": t.mutate(rn=ibis.row_number()),
"share of total in a mutate": t.mutate(share=t.v / t.v.sum()),
"lag() in a mutate": t.mutate(prev=t.v.lag()),
"random() in a mutate": t.mutate(r=ibis.random()),
"filter by membership in a second file": t.filter(t.k.isin(u.k)),
"filter against a scalar subquery": t.filter(t.v > t.v.mean()),
"limit": t.limit(3),
"distinct": t.select("k", "__row_order").distinct(),
}
print(f"parent has {n} rows\n")
print(f"{'recipe shape':40s} {'today':8s} {'relations only':15s} {'D4':7s} rows __row_order unique")
names: set[str] = set()
for label, expr in shapes.items():
build = Path(build_expr(expr, builds_dir=HOME / "builds"))
names |= {m for p in build.glob("*.yaml") for m in re.findall(r"op:\s*([A-Za-z_]+)", p.read_text())}
today = "worthy" if classify_build(build)["worthy"] else "cheap"
relations = "cheap" if relation_only_allow_list(expr) else "worthy"
proposed = "cheap" if is_cheap(expr) else "worthy"
out = expr.execute()
print(f"{label:40s} {today:8s} {relations:15s} {proposed:7s} {len(out):4d} {out['__row_order'].is_unique}")

relation_names = {c.__name__ for c in ops.Relation.__subclasses__()} | {"Read"}
not_relations = sorted(n for n in names if n not in relation_names and not hasattr(ops, n))
print(f"\nnames the classify_build regex matches that are not operations at all: {not_relations}")
print(
f"value operations it matches alongside the relations: {sorted(n for n in names if n in ('Field', 'Literal'))}"
)


if __name__ == "__main__":
main()
77 changes: 77 additions & 0 deletions scripts/spike_csv_source_identity.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,77 @@
"""ADR-008 evidence (D2, D7) and #168: is a CSV root's ordered intermediate fixed under a content hash?

``tallyman_read_csv`` reads a CSV through an ordered parquet intermediate. That file's key is
``md5(absolute path | schema | reader options)``, not content, and it is overwritten in place when the CSV's mtime
changes. The script edits a CSV, re-runs the identical recipe, and reports the build hash before and after, and what
the first version's frozen build returns afterwards, for three shapes of the root expression:

1. today: a read of the intermediate with a trailing ``order_by("original_row_order")``;
2. ADR-008 D7 without the fix for #168: a plain read of the same intermediate (no sort, so no snapshot, no digest);
3. the fix for #168: digest the CSV, clone it to a content-named path, and key the intermediate on the clone.

Raw xorq builds, no tallyman build machinery, so case 1 reports the hash only: in a real build the root is worthy
because of its Sort, and its baked snapshot goes on serving the old rows. An unchanged hash also means
``build_and_persist`` returns the existing entry, so no new version is created at all.

uv run python scripts/spike_csv_source_identity.py
"""

from __future__ import annotations

import os
import tempfile
import time
from pathlib import Path

HOME = Path(tempfile.mkdtemp(prefix="spike_csv_identity_"))
os.environ["TALLYMAN_HOME"] = str(HOME) # the ordered intermediates live under it; set before tallyman is imported

import xorq.api as xo # noqa: E402
from xorq.ibis_yaml.compiler import build_expr, load_expr # noqa: E402

from tallyman_xorq import source_identity as si # noqa: E402
from tallyman_xorq.io import _ordered_csv_parquet, tallyman_read_csv # noqa: E402

BUILDS = HOME / "builds"
CAS = HOME / "cas"
ORIGINAL = "id,amount\n1,10\n2,20\n3,30\n"
EDITED = "id,amount\n1,10\n2,999\n3,30\n4,40\n"


def today(csv: Path):
return tallyman_read_csv(str(csv))


def plain_read(csv: Path):
return xo.deferred_read_parquet(str(_ordered_csv_parquet(str(csv), None, {})))


def keyed_on_clone(csv: Path):
clone = CAS / f"{si._digest_file(csv)}{csv.suffix}"
if not clone.exists():
clone.write_bytes(csv.read_bytes()) # tallyman's ensure_cas_path makes a copy-on-write clone here
return xo.deferred_read_parquet(str(_ordered_csv_parquet(str(clone), None, {})))


def main() -> None:
CAS.mkdir()
csv = HOME / "sales.csv"
cases = (
("1. today (trailing order_by kept)", today, False),
("2. ADR-008 D7 without the #168 fix", plain_read, True),
("3. keyed on a content-named clone", keyed_on_clone, True),
)
for label, root, reread in cases:
csv.write_text(ORIGINAL)
v1 = Path(build_expr(root(csv), builds_dir=BUILDS))
time.sleep(0.05) # so the edit changes the CSV's mtime
csv.write_text(EDITED)
v2 = Path(build_expr(root(csv), builds_dir=BUILDS))
print(f"{label}: hash before the edit {v1.name}, after {v2.name}, same hash: {v1.name == v2.name}")
if reread:
rows = load_expr(v1).execute().amount.tolist()
print(f" V1's frozen build, re-read after the edit: {rows} (built from {[10, 20, 30]})")


if __name__ == "__main__":
main()
77 changes: 77 additions & 0 deletions scripts/spike_deep_page_memory.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,77 @@
"""ADR-008 evidence (Consequences): what a page request holds in memory at depth.

ADR-008 D5 serves every page as ``ORDER BY __row_order LIMIT n OFFSET k``. Its cost table measures latency. This
script measures peak process memory for the same requests, on the same file shape as ``spike_row_order_paging.py``,
and for a range request (``__row_order >= k AND __row_order < k + n``), which ADR-008 notes as a later optimization.

Each case runs in a fresh process, so its peak is its own; the first case imports xorq and does nothing, as the floor.

uv run python scripts/spike_deep_page_memory.py
"""

from __future__ import annotations

import resource
import subprocess
import sys
import tempfile
from pathlib import Path

ROW = "__row_order"
N = 3_000_000
CASES = (
"import only",
"bare limit 50",
"sorted 0",
"sorted 1000000",
"sorted 2900000",
"range 2900000",
"chart 100000",
)


def run_case(path: str, case: str) -> None:
import xorq.api as xo

kind, _, arg = case.partition(" ")
if kind != "import":
t = xo.connect().read_parquet(path)
if kind == "bare":
t.limit(50).execute()
elif kind == "sorted":
t.order_by(ROW).limit(50, offset=int(arg)).execute()
elif kind == "range":
k = int(arg)
t.filter((t[ROW] >= k) & (t[ROW] < k + 50)).order_by(ROW).execute()
elif kind == "chart":
t.order_by(ROW).limit(int(arg)).execute() # the chart pull through /api/data
peak = resource.getrusage(resource.RUSAGE_SELF).ru_maxrss
peak_mb = peak / 1e6 if sys.platform == "darwin" else peak / 1e3 # bytes on macOS, kilobytes on Linux
print(f" {case:16s} peak process memory {peak_mb:7.0f} MB")


def main() -> None:
import numpy as np
import pyarrow as pa
import pyarrow.parquet as pq

rng = np.random.default_rng(9)
cols = {"g": rng.integers(0, 200, N)}
cols |= {f"v{i}": rng.random(N) for i in range(12)}
cols[ROW] = np.arange(N)
table = pa.table(cols)
with tempfile.TemporaryDirectory() as tmp:
path = Path(tmp) / "snapshot.parquet"
pq.write_table(table, path, compression="zstd", row_group_size=1_048_576, write_page_index=True)
on_disk, as_arrow = path.stat().st_size / 1e6, table.nbytes / 1e6
print(f"file: {on_disk:.0f} MB on disk, {as_arrow:.0f} MB as Arrow, {N:,} rows x {table.num_columns} columns")
for case in CASES:
out = subprocess.run([sys.executable, __file__, str(path), case], capture_output=True, text=True)
print(out.stdout.rstrip() or out.stderr.strip()[-300:])


if __name__ == "__main__":
if len(sys.argv) == 3:
run_case(sys.argv[1], sys.argv[2])
else:
main()
Loading
Loading