From 1c0321febbb36b3f3129e999122d2c2f8edfe861 Mon Sep 17 00:00:00 2001 From: Krzysztof Macewicz Date: Thu, 24 Sep 2026 09:42:37 +0200 Subject: [PATCH] docs: in-place index builds measured with the shipped operator The premise compared the copy path with the engines' own non-blocking statements. This runs LocalCutover.index from f77bb7b against the copy path on the same tables, with a process on the map in force writing throughout. At 450 000 rows on PostgreSQL the copy path's cutover ran into the operator's 30 s watchdog and a recovery aborted it, leaving the decided index out of force; the in-place build was published in 618 ms. On ClickHouse the same table still fitted, 2.2 s inside the budget. Co-Authored-By: Claude Opus 5.5 (1M context) --- docs/in-place-index.md | 4 +- docs/qualification/in-place-index/README.md | 26 ++- .../in-place-index/after/SHA256SUMS | 5 + .../in-place-index/after/after-run.log | 9 + .../in-place-index/after/after.json | 119 ++++++++++ .../in-place-index/after/load-after.txt | 1 + .../in-place-index/after/load-before.txt | 1 + .../in-place-index/after/measure_after.py | 220 ++++++++++++++++++ 8 files changed, 383 insertions(+), 2 deletions(-) create mode 100644 docs/qualification/in-place-index/after/SHA256SUMS create mode 100644 docs/qualification/in-place-index/after/after-run.log create mode 100644 docs/qualification/in-place-index/after/after.json create mode 100644 docs/qualification/in-place-index/after/load-after.txt create mode 100644 docs/qualification/in-place-index/after/load-before.txt create mode 100644 docs/qualification/in-place-index/after/measure_after.py diff --git a/docs/in-place-index.md b/docs/in-place-index.md index a8ea8c7..7056921 100644 --- a/docs/in-place-index.md +++ b/docs/in-place-index.md @@ -11,7 +11,9 @@ under new names ([staging](staging.md) protocol 2) and a [cutover](local-cutover compares every row while the source's writes are frozen. That pause grows with the table - measured at 22.5 s for 300 000 rows on PostgreSQL, and a table of about 400 000 rows passes the cutover's 30 s budget and rolls back. The in-place build of the same index took 355 ms, and the longest gap -between two writes during it was 16 ms. Measurements and scripts: +between two writes during it was 16 ms. Measured again with the shipped operator: at 450 000 rows +on PostgreSQL the copy path's cutover ran into the 30 s watchdog and was rolled back, while the +in-place build was published in 618 ms. Measurements and scripts: [qualification/in-place-index](qualification/in-place-index/README.md). Only adding indexes can be done in place. A new key order, a partition, other tables or another diff --git a/docs/qualification/in-place-index/README.md b/docs/qualification/in-place-index/README.md index b704b30..809ab24 100644 --- a/docs/qualification/in-place-index/README.md +++ b/docs/qualification/in-place-index/README.md @@ -25,13 +25,37 @@ and `premise/premise-run.log` (100 000 and 300 000 rows). The pause through the copy is linear in the table - about 75 µs a row on PostgreSQL - because the cutover copies and compares every row while the source's writes are frozen. At about 400 000 rows -it passes the cutover's 30 s budget, and the cutover rolls back: an "add an index" decision could +on PostgreSQL it passes the cutover's 30 s budget, and the cutover rolls back: an "add an index" decision could not be carried out for a table of that size at all. The in-place build moves no row; the longest gap between two writes during it is of the order of the gaps without it. No write failed in any run. The script's usage lines were rewritten when it moved into this directory; its code is the code that produced these files. `premise/SHA256SUMS` covers every file here. +## After: the shipped operator against the copy path + +`after/measure_after.py` runs what shipped in PR #83 (`f77bb7b`): `LocalCutover.index` on a signed +`sde-index` authorization - snapshot, native build, qualification from the catalogue, watermark, +atomic publication - against the copy path on the same table, with a process on the map in force +writing a row every 5 ms throughout. Same laptop, same engines, load average 1.0-1.4 at the start +and 2.6 at the end. Raw results: `after/after.json`, `after/after-run.log`. + +| engine | rows | in place: receipt / operator wall | longest gap between writes during the build (before it) | copy path: write pause | copy path: outcome | +|---|---|---|---|---|---| +| PostgreSQL | 100 000 | 346 ms / 482 ms | 13.8 ms (9.6) | 7 174 ms | success | +| PostgreSQL | 450 000 | 618 ms / 764 ms | 80.2 ms (9.2) | interrupted by the operator's watchdog at 30 066 ms | recovery required; the resume aborted (`recovered_before_decision`), 257 ms | +| ClickHouse | 100 000 | 657 ms / 885 ms | 20.0 ms (27.3) | 5 581 ms | success | +| ClickHouse | 450 000 | 707 ms / 920 ms | 22.1 ms (21.2) | 27 843 ms | success, 2.2 s inside the 30 s budget | + +The table the copy path cannot finish is there: at 450 000 rows on PostgreSQL the cutover's pause +ran into the operator's 30 s watchdog, the source stayed frozen until a recovery aborted the +cutover, and the index the model decided was not in force - while the in-place build of the same +index was built and published in 618 ms. On ClickHouse the same table still fitted, by 2.2 s. No +write failed in any in-place build. The gap that stands out, 80 ms at 450 000 rows on PostgreSQL, +is the longest interval between two successful writes during that build; the script records only +the maximum, so how often a gap of that size occurred is not known, and this record does not +explain it. + ## Engine behaviours the rules are written against `premise/probe_engines.py` -> `premise/probe_engines.out.json`: diff --git a/docs/qualification/in-place-index/after/SHA256SUMS b/docs/qualification/in-place-index/after/SHA256SUMS new file mode 100644 index 0000000..6f478a2 --- /dev/null +++ b/docs/qualification/in-place-index/after/SHA256SUMS @@ -0,0 +1,5 @@ +087ef78ee4bc97abafa614f8930f1721f6cc8dcf554c72921f44a8af97e4aa45 measure_after.py +8f2299c2c296b68f8af70b20fff36d3196c259e4df36e422d7010304a509b921 after.json +9b54ae3bec5c50a18e5c35c9411e6ac7440968b7679e0dff916afc4e6cf2126a after-run.log +886f1840e5822f594fa1b0fd2cef0d47d1798d6299f7cf85d281f4570a680000 load-before.txt +faca8f3df3b557fc706ce971c2b413993bf184913a5e30149f5c4ebc4ec231d4 load-after.txt diff --git a/docs/qualification/in-place-index/after/after-run.log b/docs/qualification/in-place-index/after/after-run.log new file mode 100644 index 0000000..4041352 --- /dev/null +++ b/docs/qualification/in-place-index/after/after-run.log @@ -0,0 +1,9 @@ +{"path": "in_place", "dialect": "postgres", "rows": 100000, "outcome": "built", "receipt_elapsed_ms": 346, "operator_wall_ms": 482, "writer_max_gap_before_ms": 9.6, "writer_max_gap_during_ms": 13.8, "writer_errors": 0, "index_in_force": true} +{"path": "copy", "dialect": "postgres", "rows": 100000, "stage_ms": 401, "stage_outcome": "prepared", "cutover_wall_ms": 7243, "cutover_outcome": "success", "cutover_reason": "matched", "within_budget": true, "pause_ms": 7174, "recovery": null, "index_in_force": true} +{"path": "in_place", "dialect": "postgres", "rows": 450000, "outcome": "built", "receipt_elapsed_ms": 618, "operator_wall_ms": 764, "writer_max_gap_before_ms": 9.2, "writer_max_gap_during_ms": 80.2, "writer_errors": 0, "index_in_force": true} +{"path": "copy", "dialect": "postgres", "rows": 450000, "stage_ms": 392, "stage_outcome": "prepared", "cutover_wall_ms": 30397, "cutover_outcome": "recovery_required", "cutover_reason": "operator deadline interrupted an operation; reconnect dedicated adapters and inspect local status before resuming or retrying qualification", "within_budget": null, "pause_ms": null, "recovery": {"interrupted_after_ms": 30066, "resume_ms": 257, "resumed_outcome": "abort", "resumed_reason": "recovered_before_decision"}, "index_in_force": false} +{"path": "in_place", "dialect": "clickhouse", "rows": 100000, "outcome": "built", "receipt_elapsed_ms": 657, "operator_wall_ms": 885, "writer_max_gap_before_ms": 27.3, "writer_max_gap_during_ms": 20.0, "writer_errors": 0, "index_in_force": true} +{"path": "copy", "dialect": "clickhouse", "rows": 100000, "stage_ms": 641, "stage_outcome": "prepared", "cutover_wall_ms": 5782, "cutover_outcome": "success", "cutover_reason": "matched", "within_budget": true, "pause_ms": 5581, "recovery": null, "index_in_force": true} +{"path": "in_place", "dialect": "clickhouse", "rows": 450000, "outcome": "built", "receipt_elapsed_ms": 707, "operator_wall_ms": 920, "writer_max_gap_before_ms": 21.2, "writer_max_gap_during_ms": 22.1, "writer_errors": 0, "index_in_force": true} +{"path": "copy", "dialect": "clickhouse", "rows": 450000, "stage_ms": 652, "stage_outcome": "prepared", "cutover_wall_ms": 28054, "cutover_outcome": "success", "cutover_reason": "matched", "within_budget": true, "pause_ms": 27843, "recovery": null, "index_in_force": true} +EXIT=0 diff --git a/docs/qualification/in-place-index/after/after.json b/docs/qualification/in-place-index/after/after.json new file mode 100644 index 0000000..c3f70af --- /dev/null +++ b/docs/qualification/in-place-index/after/after.json @@ -0,0 +1,119 @@ +{ + "machine": { + "node": "km-ThinkPad-E570", + "processor": "x86_64", + "python": "3.12.3" + }, + "sde_version": "0.1.0.dev0", + "results": [ + { + "path": "in_place", + "dialect": "postgres", + "rows": 100000, + "outcome": "built", + "receipt_elapsed_ms": 346, + "operator_wall_ms": 482, + "writer_max_gap_before_ms": 9.6, + "writer_max_gap_during_ms": 13.8, + "writer_errors": 0, + "index_in_force": true + }, + { + "path": "copy", + "dialect": "postgres", + "rows": 100000, + "stage_ms": 401, + "stage_outcome": "prepared", + "cutover_wall_ms": 7243, + "cutover_outcome": "success", + "cutover_reason": "matched", + "within_budget": true, + "pause_ms": 7174, + "recovery": null, + "index_in_force": true + }, + { + "path": "in_place", + "dialect": "postgres", + "rows": 450000, + "outcome": "built", + "receipt_elapsed_ms": 618, + "operator_wall_ms": 764, + "writer_max_gap_before_ms": 9.2, + "writer_max_gap_during_ms": 80.2, + "writer_errors": 0, + "index_in_force": true + }, + { + "path": "copy", + "dialect": "postgres", + "rows": 450000, + "stage_ms": 392, + "stage_outcome": "prepared", + "cutover_wall_ms": 30397, + "cutover_outcome": "recovery_required", + "cutover_reason": "operator deadline interrupted an operation; reconnect dedicated adapters and inspect local status before resuming or retrying qualification", + "within_budget": null, + "pause_ms": null, + "recovery": { + "interrupted_after_ms": 30066, + "resume_ms": 257, + "resumed_outcome": "abort", + "resumed_reason": "recovered_before_decision" + }, + "index_in_force": false + }, + { + "path": "in_place", + "dialect": "clickhouse", + "rows": 100000, + "outcome": "built", + "receipt_elapsed_ms": 657, + "operator_wall_ms": 885, + "writer_max_gap_before_ms": 27.3, + "writer_max_gap_during_ms": 20.0, + "writer_errors": 0, + "index_in_force": true + }, + { + "path": "copy", + "dialect": "clickhouse", + "rows": 100000, + "stage_ms": 641, + "stage_outcome": "prepared", + "cutover_wall_ms": 5782, + "cutover_outcome": "success", + "cutover_reason": "matched", + "within_budget": true, + "pause_ms": 5581, + "recovery": null, + "index_in_force": true + }, + { + "path": "in_place", + "dialect": "clickhouse", + "rows": 450000, + "outcome": "built", + "receipt_elapsed_ms": 707, + "operator_wall_ms": 920, + "writer_max_gap_before_ms": 21.2, + "writer_max_gap_during_ms": 22.1, + "writer_errors": 0, + "index_in_force": true + }, + { + "path": "copy", + "dialect": "clickhouse", + "rows": 450000, + "stage_ms": 652, + "stage_outcome": "prepared", + "cutover_wall_ms": 28054, + "cutover_outcome": "success", + "cutover_reason": "matched", + "within_budget": true, + "pause_ms": 27843, + "recovery": null, + "index_in_force": true + } + ] +} diff --git a/docs/qualification/in-place-index/after/load-after.txt b/docs/qualification/in-place-index/after/load-after.txt new file mode 100644 index 0000000..f246b32 --- /dev/null +++ b/docs/qualification/in-place-index/after/load-after.txt @@ -0,0 +1 @@ + 09:41:38 up 1 day, 13:32, 1 user, load average: 2,60, 1,60, 1,53 diff --git a/docs/qualification/in-place-index/after/load-before.txt b/docs/qualification/in-place-index/after/load-before.txt new file mode 100644 index 0000000..d986de2 --- /dev/null +++ b/docs/qualification/in-place-index/after/load-before.txt @@ -0,0 +1 @@ + 09:39:04 up 1 day, 13:29, 1 user, load average: 1,01, 1,24, 1,44 diff --git a/docs/qualification/in-place-index/after/measure_after.py b/docs/qualification/in-place-index/after/measure_after.py new file mode 100644 index 0000000..c0083e3 --- /dev/null +++ b/docs/qualification/in-place-index/after/measure_after.py @@ -0,0 +1,220 @@ +"""In-place index builds measured after they exist: the real operator against the copy path. + +The premise (``../premise``) compared the copy path with the engines' own non-blocking statements. +This runs the shipped operator instead - ``LocalCutover.index`` on a signed ``sde-index`` +authorization, with its snapshot, native build, qualification from the catalogue, watermark and +atomic publication - on the same table sizes and against the same copy path, including a table the +copy path cannot finish inside the 30 s cutover budget. + +For each engine and size: + +- copy path: the staging live test's fixture with a designed copy in the source's engine (staging + protocol 2), N rows, ``stage`` then the matching cutover; the pause is the cutover receipt's + ``elapsed_ms`` and its outcome - past the budget the cutover aborts and the index is not there; +- in place: the in-place live test's fixture, the same N rows, a process on the map in force + writing a row every 5 ms on its own connection, and ``LocalCutover.index``; the result is the + receipt, the operator's wall time, and the writer's longest gap between two successful writes + during the build against its longest gap before it. + +It imports the live-test fixtures, so it runs from ``python/tests`` with both engine DSNs set: + + cd python/tests + ../.venv/bin/python ../../docs/qualification/in-place-index/after/measure_after.py \\ + --rows 100000 450000 --out after.json +""" + +from __future__ import annotations + +import argparse +import json +import platform +import sys +import threading +import time +from pathlib import Path +from tempfile import TemporaryDirectory +from typing import Any + +sys.path.insert(0, str(Path.cwd())) + +import test_index_operator_live as in_place_fixture +from test_staging_operator_live import cutover, initial + +import sde +from sde.engines.clickhouse import ClickHouseEngine +from sde.engines.postgres import PostgresEngine +from sde.generation import EPOCH_COLUMN +from sde.local_cutover import CutoverRecoveryRequired, LocalCutover + +CHUNK = 20000 + + +def seed(engine: Any, table: str, first: int, count: int) -> None: + for start in range(first, first + count, CHUNK): + rows = [ + {"id": i, "value": i % 100000, EPOCH_COLUMN: 1} + for i in range(start, min(first + count, start + CHUNK)) + ] + engine.copy_in(table, rows) + + +def index_for(dialect: str, name: str) -> dict[str, Any]: + if dialect == "postgres": + return {"entity": "Event", "name": name, "columns": ["value"]} + return { + "entity": "Event", + "name": name, + "columns": ["value"], + "method": "minmax", + "granularity": 4, + } + + +def copy_path(dialect: str, rows: int) -> dict[str, Any]: + with ( + TemporaryDirectory() as tmp, + initial( + dialect, Path(tmp), target=dialect, indexes=[index_for(dialect, "sde_i_after_000001")] + ) as (operator, stage, _old, roles, model, public, signed, _m, _s, _t), + ): + seed(roles[dialect].operator, "initial_events", 2, rows - 1) + started = time.monotonic() + staged = operator.stage(stage).as_record() + stage_s = time.monotonic() - started + started = time.monotonic() + recovery: dict[str, Any] | None = None + try: + receipt = operator.execute(cutover(stage, signed, model, public)).as_record() + outcome, reason, elapsed = receipt["outcome"], receipt["reason"], receipt["elapsed_ms"] + within = receipt["within_budget"] + except CutoverRecoveryRequired as exc: + # The operator's watchdog ended the pause: its connections are closed and the source is + # still frozen. Fresh connections and a resume - without a decision, recovery aborts + # and reopens the source - are what a customer would do next; both are recorded. + outcome, reason, elapsed, within = "recovery_required", str(exc), None, None + wall = time.monotonic() - started + for name, role in roles.items(): + role.operator.close() + role.operator.connect() + if name == "postgres": + role.operator._cx.execute('SET search_path TO "' + role.namespace + '"') + role.runtime.close() + role.runtime.connect() + resumed_at = time.monotonic() + recovered = ( + LocalCutover( + Path(tmp), + model=model, + project_id="1" * 32, + public_key=public, + operators={name: role.operator for name, role in roles.items()}, + runtime={name: [role.runtime] for name, role in roles.items()}, + ) + .resume() + .as_record() + ) + recovery = { + "interrupted_after_ms": round(wall * 1000), + "resume_ms": round((time.monotonic() - resumed_at) * 1000), + "resumed_outcome": recovered["outcome"], + "resumed_reason": recovered["reason"], + } + return { + "path": "copy", + "dialect": dialect, + "rows": rows, + "stage_ms": round(stage_s * 1000), + "stage_outcome": staged["outcome"], + "cutover_wall_ms": round((time.monotonic() - started) * 1000), + "cutover_outcome": outcome, + "cutover_reason": reason, + "within_budget": within, + "pause_ms": elapsed, + "recovery": recovery, + "index_in_force": outcome == "success", + } + + +def in_place(dialect: str, rows: int) -> dict[str, Any]: + with TemporaryDirectory() as tmp, in_place_fixture.initial(dialect, Path(tmp)) as build: + seed(build.role.operator, in_place_fixture.TABLE, 2, rows - 1) + engines = { + name: (PostgresEngine if name == "postgres" else ClickHouseEngine)(role.runtime._dsn) + for name, role in build.roles.items() + } + for engine in engines.values(): + engine.connect() + writer = sde.Session(build.model, build.plan.current, engines, project_id="1" * 32) + gaps = {"before": 0.0, "during": 0.0} + phase = {"name": "before"} + stop = threading.Event() + errors: list[str] = [] + + def write() -> None: + last = time.monotonic() + identity = rows + 1 + while not stop.is_set(): + try: + writer.save("Event", {"id": identity, "value": 1}) + now = time.monotonic() + gaps[phase["name"]] = max(gaps[phase["name"]], now - last) + last = now + identity += 1 + except Exception as exc: + errors.append(type(exc).__name__) + time.sleep(0.005) + + thread = threading.Thread(target=write) + thread.start() + time.sleep(2.0) + phase["name"] = "during" + started = time.monotonic() + receipt = build.operator.index(build.plan).as_record() + wall_s = time.monotonic() - started + time.sleep(0.5) + stop.set() + thread.join() + for engine in engines.values(): + engine.close() + return { + "path": "in_place", + "dialect": dialect, + "rows": rows, + "outcome": receipt["outcome"], + "receipt_elapsed_ms": receipt["elapsed_ms"], + "operator_wall_ms": round(wall_s * 1000), + "writer_max_gap_before_ms": round(gaps["before"] * 1000, 1), + "writer_max_gap_during_ms": round(gaps["during"] * 1000, 1), + "writer_errors": len(errors), + "index_in_force": receipt["outcome"] == "built", + } + + +def main() -> int: + parser = argparse.ArgumentParser() + parser.add_argument("--rows", type=int, nargs="+", required=True) + parser.add_argument("--dialects", nargs="+", default=["postgres", "clickhouse"]) + parser.add_argument("--out", type=Path, required=True) + args = parser.parse_args() + results = [] + for dialect in args.dialects: + for rows in args.rows: + for path in (in_place, copy_path): + result = path(dialect, rows) + results.append(result) + print(json.dumps(result), flush=True) + report = { + "machine": { + "node": platform.node(), + "processor": platform.processor(), + "python": platform.python_version(), + }, + "sde_version": sde.__version__, + "results": results, + } + args.out.write_text(json.dumps(report, indent=2) + "\n") + return 0 + + +if __name__ == "__main__": + raise SystemExit(main())