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
5 changes: 3 additions & 2 deletions docs/caching.md
Original file line number Diff line number Diff line change
Expand Up @@ -169,7 +169,8 @@ A recipe never reads a source file directly. `read_project_file` and
`tallyman_read_csv` take the source's content digest (md5, memoized on the
file's stat). In the default `cas` mode they also clone the source
copy-on-write to `<project>/data/.cas/<digest><suffix>`, the **clone**. Then
polars writes an **ordered copy**: `<compute_cache>/ordered_sources/<key>.parquet`,
pyarrow (a parquet source, whose types it keeps) or polars (a CSV, which it
parses) writes an **ordered copy**: `<compute_cache>/ordered_sources/<key>.parquet`,
holding the source's rows in file order plus a last column `__row_order`, in
row groups of 122,880 rows. The recipe reads the copy. `<key>` is an md5 of the
digest and the reader options (for a CSV, the schema and the `scan_csv`
Expand Down Expand Up @@ -223,7 +224,7 @@ its own snapshot, is on disk before anything executes:
| File | Written by | If it is missing |
|---|---|---|
| Snapshot, `compute_cache/result_cache/<hash>.parquet` | `materialize` | re-run the entry's build, and verify the digest |
| Ordered copy, `compute_cache/ordered_sources/<key>.parquet` | polars, at ingest | re-run ingest on the clone with the recorded reader options, and check it |
| Ordered copy, `compute_cache/ordered_sources/<key>.parquet` | pyarrow or polars, at ingest | re-run ingest on the clone with the recorded reader options, and check it |
| Clone, `data/.cas/<digest><suffix>` | `ensure_cas_path` | copy the live source again, but only while its bytes still hash to the digest |

A healed snapshot is checked against the recorded `result_digest`. A mismatch
Expand Down
5 changes: 3 additions & 2 deletions docs/system-contract.md
Original file line number Diff line number Diff line change
Expand Up @@ -193,8 +193,9 @@ input forever; the live file is merely where the *next* build will look.

A recipe never reads the source or its clone directly. Ingest writes an
**ordered copy** of the clone under the project's `compute_cache/ordered_sources/`:
polars reads the clone in file order and writes a parquet file with the same
columns and one more at the end, `__row_order`, `0..N-1`. The copy is named by
pyarrow (a parquet source, whose types it keeps) or polars (a CSV) reads the
clone in file order and writes a parquet file with the same columns and one more
at the end, `__row_order`, `0..N-1`. The copy is named by
the source's digest and the reader options (the schema and `scan_csv` options
of a CSV), and the recipe reads that file. So a recipe's reads are all files
tallyman wrote, each carrying the column that pages sort by (Part 2, "Row
Expand Down
6 changes: 6 additions & 0 deletions plans/ADR-008-row-order-of-reads.md
Original file line number Diff line number Diff line change
Expand Up @@ -702,6 +702,12 @@ that does not exist yet fails on import, and that counts as red. Paddy,
- **Open question 1 and 2** were implemented as the uniform rule: every parquet source gets an ordered copy, and the
column is `__row_order`. **Open question 5** is answered by ADR-007 D13: the manifest's `ordered_copies` records the
copy's content digest, and a re-created copy is checked against it (a mismatch is loud and served).
- **D2, who writes the copy.** pyarrow copies a parquet source (#197): it streams the source in file order and writes
the copy in the snapshot's parquet settings with 122,880-row groups, so the copy's schema is the source's plus
`__row_order`. polars wrote it at first and did not keep every type: a `date64` came back as a timestamp, a map as
a list of structs, `time32` and `time64` as `time64[ns]`, and a `decimal256` made it panic. polars still parses a
CSV, and a panic there is a `BuildError`. The first ingest of a CSV passes the reader options as the caller gave
them; the JSON form the manifest records is only replayed to make a deleted copy again (#198).
- **A CSV that already has a `__row_order` column** has it overwritten, like a parquet source; the previous check that
the column was a canonical `0..N-1` sequence is gone, and `original_row_order` is ordinary data.

Expand Down
7 changes: 7 additions & 0 deletions plans/ADR-009-digest-stability.md
Original file line number Diff line number Diff line change
Expand Up @@ -444,6 +444,13 @@ that does not exist yet fails on import, and that counts as red. Paddy,
- **D3.** `SNAPSHOT_FORMAT_VERSION = 1` stands for the snapshot row-group size (1,048,576), the materialization
connection's batch size (8,192) and the ordered-copy row-group size (122,880). It is recorded in the manifest with
the xorq, xorq-datafusion and pyarrow versions.
- **D3, the ordered copies.** Since #197 pyarrow writes the copy of a parquet source, in this writer's parquet
settings with 122,880-row groups; polars writes only a CSV's. The format version stayed at 1: the row groups are
where they were, and polars also wrote a page index. On the 1.5M-row test source the pyarrow copy has the same
content digest as the polars one, and on the single-partition connection the `SUM` and `AVG` of its float column
come out bit for bit the same, ungrouped, filtered on the float, and over four ranges of the sorted `id` that the
page index prunes. What changed is the columns polars did not keep (`date64`, maps, `time32`, `time64`), and a
copy of those made again is reported by its recorded content digest (`unfaithful_ordered_copy`).
- **D4.** The heal record says the engine changed, naming the versions, when any recorded version or the format
differs from today's, and otherwise keeps the structural (#88) and execution (#83) attribution.

Expand Down
8 changes: 6 additions & 2 deletions src/tallyman_xorq/io.py
Original file line number Diff line number Diff line change
Expand Up @@ -471,7 +471,11 @@ def tallyman_read_csv(path: str, schema=None, project: str | None = None, **kwar
**kwargs: Forwarded to ``polars.scan_csv`` — reader options such as
``separator``, ``skip_rows``, ``null_values``, ``quote_char``,
``has_header``, ``encoding``. They participate in the copy's key,
so changing one re-ingests. ``infer_schema_length`` and
so changing one re-ingests. The first read passes them as given;
the manifest records them as JSON, which keeps only the repr of a
function (``with_column_names``), so a deleted copy of such a read
cannot be made again and its entry has to be rebuilt (#198).
``infer_schema_length`` and
``schema_overrides`` are managed internally (see ``_RESERVED_SCAN_KWARGS``)
and rejected — they would collide with the values every internal
``scan_csv`` call already sets.
Expand Down Expand Up @@ -501,7 +505,7 @@ def tallyman_read_csv(path: str, schema=None, project: str | None = None, **kwar
if si.mode() != "off":
si.note_source(rel, digest)
source = si.ensure_cas_path(proj, src, digest) if si.mode() == "cas" else src
copy = oc.ensure_ordered_copy(proj, source, digest=digest, rel=rel, reader=reader)
copy = oc.ensure_ordered_copy(proj, source, digest=digest, rel=rel, reader=reader, csv_args=(schema, kwargs))
return deferred_read_parquet(str(copy))


Expand Down
152 changes: 118 additions & 34 deletions src/tallyman_xorq/ordered_copy.py
Original file line number Diff line number Diff line change
@@ -1,8 +1,9 @@
"""Ordered copies of sources (ADR-008 D2, ADR-007 D13): the one way a file enters a recipe.

A source (a parquet file or a CSV under the project) is read through its content-addressed clone (``data/.cas``), and
polars writes a parquet copy of it, in file order, with one more column at the end: ``__row_order``, ``0..N-1``. That
copy is what a recipe reads. Nothing reads the source or the clone directly, so:
a parquet copy of it is written, in file order, with one more column at the end: ``__row_order``, ``0..N-1``. pyarrow
copies a parquet source, so the copy keeps the source's types (#197); polars parses a CSV. That copy is what a recipe
reads. Nothing reads the source or the clone directly, so:

- every file tallyman reads carries the column that pages sort by;
- the copy is keyed by the source's content digest and the reader options, so editing a source and running the same
Expand All @@ -25,12 +26,18 @@
import uuid
from pathlib import Path

import numpy as np
import pyarrow as pa
import pyarrow.parquet as pq

from tallyman_xorq.row_order import ROW_ORDER

perf_log = logging.getLogger("tallyman.perf")

# Pinned polars parquet-write settings for an ordered copy. Held constant so the layout, and therefore any float total
# computed straight from a source, is reproducible. Changing one is a corpus rebuild (SNAPSHOT_FORMAT_VERSION).
# The row-group size of an ordered copy. Held constant so the layout, and therefore any float total computed straight
# from a source, is reproducible. Changing it is a corpus rebuild (SNAPSHOT_FORMAT_VERSION). A parquet source's copy is
# otherwise written in the snapshot's parquet settings (``materialize._PARQUET_OPTIONS``); polars writes a CSV's with
# the settings below.
ORDERED_COPY_ROW_GROUP_ROWS = 122_880
_WRITE = {
"compression": "zstd",
Expand Down Expand Up @@ -142,65 +149,142 @@ def end_collect(token: contextvars.Token) -> dict[str, dict]:
# ---------------------------------------------------------------------------


def _write_parquet_copy(src: Path, dest: Path) -> None:
import polars as pl
def _row_groups(batches, rows: int):
"""The rows of *batches*, in order, as tables of *rows* rows each; the last one may be shorter."""
pending: list[pa.RecordBatch] = []
pending_rows = 0
for batch in batches:
if not batch.num_rows:
continue
pending.append(batch)
pending_rows += batch.num_rows
while pending_rows >= rows:
table = pa.Table.from_batches(pending)
yield table.slice(0, rows)
tail = table.slice(rows)
pending, pending_rows = tail.to_batches(), tail.num_rows
if pending_rows:
yield pa.Table.from_batches(pending)

lf = pl.scan_parquet(str(src))
columns = [c for c in lf.collect_schema().names() if c != ROW_ORDER] # an existing __row_order is overwritten
lf.select(columns).with_row_index(ROW_ORDER).select([*columns, pl.col(ROW_ORDER).cast(pl.Int64)]).sink_parquet(
str(dest), **_WRITE
)


def _write_copy(src: Path, reader: dict, dest: Path) -> None:
def _write_parquet_copy(src: Path, dest: Path) -> None:
"""Copy the parquet file *src* to *dest* with pyarrow, in file order, numbering the rows in a last ``__row_order``.

pyarrow reads and writes every type a parquet file can hold, so the copy's schema is the source's plus
``__row_order``, and a recipe sees the types a direct read of the source gives. polars, which wrote the copy
before, turned a ``date64`` into a timestamp and a map into a list of structs, and panicked on a ``decimal256``
(#197). An existing ``__row_order`` is dropped and written again, last. The rows are regrouped into row groups of
``ORDERED_COPY_ROW_GROUP_ROWS``, whatever the source's row groups are, and each is combined into contiguous arrays,
as the snapshot writer does (``materialize._stream_to_parquet``). Memory is bounded by one row group.
"""
from tallyman_xorq.materialize import _PARQUET_OPTIONS

with pq.ParquetFile(src) as source:
kept = [f for f in source.schema_arrow if f.name != ROW_ORDER]
schema = pa.schema([*kept, pa.field(ROW_ORDER, pa.int64())])
written = 0
with pq.ParquetWriter(dest, schema, **_PARQUET_OPTIONS) as writer:
batches = source.iter_batches(columns=[f.name for f in kept])
for table in _row_groups(batches, ORDERED_COPY_ROW_GROUP_ROWS):
numbers = pa.array(np.arange(written, written + table.num_rows, dtype=np.int64))
numbered = table.append_column(schema.field(ROW_ORDER), numbers)
writer.write_table(numbered.combine_chunks(), row_group_size=ORDERED_COPY_ROW_GROUP_ROWS)
written += table.num_rows


def _write_copy(src: Path, reader: dict, dest: Path, *, rel: str, csv_args: tuple | None = None) -> None:
"""Write the ordered copy of *src*, the source *rel*, to *dest*.

A CSV is parsed with *csv_args*, the ``(schema, scan_kwargs)`` that ``tallyman_read_csv`` was called with, when they
are given: that is the first write. A re-creation has only the manifest's record of them, which went through JSON,
and JSON turns a function into its repr (#198); ``recreate_ordered_copy`` refuses a record marked
``lossless: False``.
"""
if reader["kind"] == "parquet":
_write_parquet_copy(src, dest)
return
import polars as pl

from tallyman_xorq.io import _materialize_ordered

_materialize_ordered(src, _spec_from_json(reader["schema"]), dict(reader["scan_kwargs"]), dest)
if csv_args is None:
csv_args = (_spec_from_json(reader["schema"]), dict(reader["scan_kwargs"]))
schema, scan_kwargs = csv_args
try:
_materialize_ordered(src, schema, scan_kwargs, dest)
except pl.exceptions.PanicException as exc:
# A PanicException is a BaseException, so no `except Exception` on the build or MCP path would catch it (#197).
from tallyman_xorq.build import BuildError

raise BuildError(f"tallyman_read_csv: polars panicked reading {rel!r}: {exc}") from exc


def _digest_sidecar(path: Path) -> Path:
return path.with_suffix(".digest")


def _write_sidecar(path: Path, digest: str) -> None:
"""Write the digest sidecar of the copy *path* to a unique temp name and replace, so no reader sees it half done."""
sidecar = _digest_sidecar(path)
tmp = sidecar.with_name(f"{sidecar.name}.{uuid.uuid4().hex}.tmp")
try:
tmp.write_text(digest)
os.replace(tmp, sidecar)
finally:
tmp.unlink(missing_ok=True)


def _content_digest_of(path: Path) -> str:
"""The copy's content digest: the sidecar written with it, or a read-back when the sidecar is gone."""
"""The copy's content digest: the sidecar written with it, or a read-back when the sidecar is missing or empty.

An empty sidecar is what an in-place write cut off between its truncate and its write leaves behind, and taking it
would record an empty digest (#211).
"""
from tallyman_xorq.digest import content_digest

sidecar = _digest_sidecar(path)
try:
return sidecar.read_text().strip()
recorded = _digest_sidecar(path).read_text().strip()
except OSError:
digest = content_digest(path)
try:
sidecar.write_text(digest)
except OSError:
pass
return digest
recorded = ""
if recorded:
return recorded
digest = content_digest(path)
try:
_write_sidecar(path, digest)
except OSError:
pass
return digest


def _write_atomically(src: Path, reader: dict, target: Path, *, rel: str, csv_args: tuple | None = None) -> str:
"""Write the copy to a unique temp name beside its destination, then its digest sidecar, then replace the copy into
place; return the content digest.

def _write_atomically(src: Path, reader: dict, target: Path) -> str:
"""Write the copy to a unique temp name beside its destination, replace, and return its content digest."""
Both files arrive by an atomic replace, the sidecar first, so a reader that finds the copy finds its whole digest
beside it, never an empty one or one left by an earlier copy (#211).
"""
from tallyman_xorq.digest import content_digest

target.parent.mkdir(parents=True, exist_ok=True)
tmp = target.with_name(f"{target.stem}.{uuid.uuid4().hex}.tmp")
try:
_write_copy(src, reader, tmp)
_write_copy(src, reader, tmp, rel=rel, csv_args=csv_args)
digest = content_digest(tmp)
_write_sidecar(target, digest)
os.replace(tmp, target)
finally:
tmp.unlink(missing_ok=True)
digest = content_digest(target)
_digest_sidecar(target).write_text(digest)
return digest


def ensure_ordered_copy(project: str, source: Path, *, digest: str, rel: str, reader: dict) -> Path:
def ensure_ordered_copy(
project: str, source: Path, *, digest: str, rel: str, reader: dict, csv_args: tuple | None = None
) -> Path:
"""The ordered copy of *source* (whose content digest is *digest*), written if it is not on disk yet.

*source* is the file polars reads: the clone in cas mode, the live file otherwise. The copy is recorded for the
*source* is the file that is read: the clone in cas mode, the live file otherwise. For a CSV, *csv_args* are the
``(schema, scan_kwargs)`` its caller passed, and the first write parses the CSV with them as they are; *reader*
holds their JSON form, which names the copy and is what a re-creation replays (#198). The copy is recorded for the
build in progress, so its manifest can make it again.
"""
from tallyman_core.catalog_state import project_lock
Expand All @@ -210,7 +294,7 @@ def ensure_ordered_copy(project: str, source: Path, *, digest: str, rel: str, re
if not target.exists():
with project_lock(project):
if not target.exists(): # a peer may have written it while we waited
_write_atomically(source, reader, target)
_write_atomically(source, reader, target, rel=rel, csv_args=csv_args)
note_ordered_copy(
key,
{
Expand Down Expand Up @@ -314,12 +398,12 @@ def recreate_ordered_copy(project: str, owner_hash: str, path: Path) -> None:
if path.exists():
return
src = _clone_or_live(project, record)
digest = _write_atomically(src, reader, path)
digest = _write_atomically(src, reader, path, rel=record["source"])
if digest != record["content_digest"]:
message = (
f"the ordered copy of {record['source']!r} was made again with content digest {digest}, not the "
f"{record['content_digest']} recorded when it was first written (a change in polars, or a source that is "
"not what it was)"
f"{record['content_digest']} recorded when it was first written (a change in pyarrow or polars, which "
"write the copies, or a source that is not what it was)"
)
perf_log.warning("recreate_ordered_copy %s: %s", key, message)
try:
Expand Down
17 changes: 17 additions & 0 deletions tests/test_materialize.py
Original file line number Diff line number Diff line change
Expand Up @@ -456,6 +456,23 @@ def test_the_ordered_copy_lives_under_the_compute_cache_and_is_recorded(project,
assert record["content_digest"] == snapshot_file_digest(copies[0])


def test_an_empty_digest_sidecar_is_recomputed_not_recorded(project, orders_parquet, monkeypatch):
"""#211: the ``.digest`` file beside a copy caches its content digest. One left empty (a write cut off between its
truncate and its write) is treated as missing: the next build that reads the source takes the digest from the copy
and records that, not ``""``, which would make a later re-creation look unfaithful."""
monkeypatch.setenv("TALLYMAN_PROJECT", project)
_hash(catalog_create("orders", _root_code(project)))
[copy] = _ordered_copies(project)
sidecar = copy.with_suffix(".digest")
sidecar.write_text("")

h = build_and_persist(project, _agg_code(project)).content_hash # a second recipe over the same source

digest = snapshot_file_digest(copy)
assert read_manifest(entry_dir(project, h)).ordered_copies[copy.stem]["content_digest"] == digest
assert sidecar.read_text() == digest


def test_a_deleted_ordered_copy_is_made_again_from_the_clone_and_checked(project, orders_parquet, monkeypatch):
"""ADR-007 D13 and D5: opening a root entry re-creates its ordered copy from the clone."""
monkeypatch.setenv("TALLYMAN_PROJECT", project)
Expand Down
Loading
Loading