Skip to content
Merged
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: 5 additions & 0 deletions .gitignore
Original file line number Diff line number Diff line change
Expand Up @@ -31,3 +31,8 @@ docs/architecture-performance-map.json
docs/disparity-report.json
docs/disparity-history.json
docs/scaled-ci-history.json
docs/estimator-comparison.json

# Raw ESTIMATOR| capture from the generated wide benchmark; the joined result
# is benchmarks/estimator_comparison.json.
benchmarks/estimator_flow_raw.txt
58 changes: 58 additions & 0 deletions benchmarks/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -225,6 +225,64 @@ python benchmarks/summarize_scaled_ci.py --baseline benchmarks/scaled_flow_basel

Each run directory needs `scaled_flow.json` beside `scaled_comparison.json`. Build it from CI artifacts; a developer machine's numbers would make a fast laptop the standard CI has to meet.

## Every estimator

The canonical benchmark races 12 estimators. `lib/scikit` exports 203, and a
claim about Flow against scikit-learn says little while the other 191 are
unmeasured. Four pieces cover the rest, all driven from one registry so the two
sides cannot drift apart.

[`estimator_coverage.py`](estimator_coverage.py) parses every exported `*_fit`
out of `lib/scikit`, resolves its arguments from a table keyed by parameter
name, finds the scikit-learn class by name, and sorts each estimator into a
bucket. Nothing is dropped silently; every entry that is not raced carries its
reason.

| bucket | meaning |
| --- | --- |
| `runnable` | arguments resolved and a scikit-learn counterpart exists |
| `different_shape` | takes a pipeline, a vectorizer input or a list of fitted models first, so it is not an estimator over a feature matrix |
| `flow_only` | Flow implements it and scikit-learn has no equivalent |
| `simplified` | the implementation's own comments call it a simplified stand-in |
| `blocked` | the signature is not resolved yet, with the missing parameter named |

The `simplified` bucket is detected from the source rather than listed, so it
stays true as the implementations are filled in. It matters: `spectral_biclustering`
thresholds row and column means where scikit-learn does an SVD and k-means, and
came out at 25106x. Three of the four it catches would otherwise have been the
widest wins in the whole matrix.

[`generate_estimator_bench.py`](generate_estimator_bench.py) emits the Flow
timing blocks from that registry, split across several files because Flow issue
#469 miscompiles some functions once a program grows past a certain size. Each
block prints one line and flushes it. Without the flush, one estimator trapping
takes the whole file's buffered output with it and the run looks empty rather
than partial, which is how a RANSAC crash first presented.

[`bench_estimators_sklearn.py`](bench_estimators_sklearn.py) times the
scikit-learn side of the same registry, and
[`compare_estimators.py`](compare_estimators.py) joins them.

```
python benchmarks/estimator_coverage.py
python benchmarks/generate_estimator_bench.py
for f in benchmarks/generated/bench_estimators_*.flow; do flow run "$f"; done | tee /tmp/flow.txt
python benchmarks/bench_estimators_sklearn.py
python benchmarks/compare_estimators.py /tmp/flow.txt
```

Read the result for what it is. These rows have no parity contract, no declared
tolerances and no disparity report: each library runs its own defaults over the
same data. That answers whether a Flow implementation is in the same
performance league and says nothing about whether it computes the same thing.
The per-estimator numerical contract remains the canonical benchmark's job, and
`estimator_comparison.json` says so in its own `contract` field.

The breadth is worth having for correctness as much as for speed. The first run
of it found a heap-buffer-overflow in `_solve_lstsq_qr`, which assumed a design
has at least as many rows as columns; RANSAC reaches it by fitting 5 sampled
rows against the 10 features of diabetes.

## Pages publication

[`publish_headline_v2.py`](publish_headline_v2.py) validates the committed canonical benchmark and architecture map, then copies the JSON artifacts into `docs/` for the static site. The public benchmark and architecture pages render those artifacts directly instead of embedding hand-maintained timing claims.
164 changes: 164 additions & 0 deletions benchmarks/bench_estimators_sklearn.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,164 @@
#!/usr/bin/env python3
"""Time the scikit-learn counterpart of every runnable registry entry.

Reads `estimator_coverage.json` so the two sides cannot drift apart: an
estimator is raced here exactly when the Flow generator emitted a block for it.

Each estimator is constructed with its own defaults rather than a translation
of the Flow hyperparameters. The Flow side passes a fixed set of small values,
and matching them one by one across 172 estimators would be a second source of
error without making the comparison more honest: this measures each library
doing its own default thing on the same data, and the per-estimator parameter
contract stays the canonical benchmark's job.
"""
from __future__ import annotations

import argparse
import json
import time
import warnings
from pathlib import Path

import numpy as np
from sklearn.datasets import load_diabetes, load_iris
from sklearn.utils import all_estimators

# IterativeImputer is behind an experimental flag.
from sklearn.experimental import enable_iterative_imputer # noqa: F401

ROOT = Path(__file__).resolve().parents[1]
REGISTRY = ROOT / "benchmarks" / "estimator_coverage.json"

# Meta-estimators need a base estimator, and a few refuse their own defaults on
# this data. A shallow tree keeps the wrapper's own overhead visible rather than
# burying it under the base estimator's work.
def _constructors() -> dict:
from sklearn.linear_model import Ridge
from sklearn.tree import DecisionTreeClassifier, DecisionTreeRegressor
import numpy as np

tree_c = lambda: DecisionTreeClassifier(max_depth=3)
tree_r = lambda: DecisionTreeRegressor(max_depth=3)
return {
"ClassifierChain": lambda c: c(tree_c()),
"MultiOutputClassifier": lambda c: c(tree_c()),
"MultiOutputRegressor": lambda c: c(tree_r()),
"OneVsOneClassifier": lambda c: c(tree_c()),
"OneVsRestClassifier": lambda c: c(tree_c()),
"OutputCodeClassifier": lambda c: c(tree_c()),
"RegressorChain": lambda c: c(tree_r()),
"RFE": lambda c: c(tree_c()),
"RFECV": lambda c: c(tree_c()),
"SequentialFeatureSelector": lambda c: c(tree_c(), n_features_to_select=2),
"StackingClassifier": lambda c: c(estimators=[("t", tree_c())]),
"StackingRegressor": lambda c: c(estimators=[("r", Ridge())]),
"VotingClassifier": lambda c: c(estimators=[("t", tree_c())]),
"VotingRegressor": lambda c: c(estimators=[("r", Ridge())]),
"SparseCoder": lambda c: c(dictionary=np.eye(4)),
# nu=0.5 is infeasible for this class balance.
"NuSVC": lambda c: c(nu=0.1),
}


CLASSIFICATION = {"classification"}


def dataset_kind(entry: dict) -> str:
names = [p["name"] for p in entry["fit"]["parameters"]]
types = {p["name"]: p["type"] for p in entry["fit"]["parameters"]}
if "y" not in names and "Y" not in names:
return "unsupervised"
if "Y" in names:
return "multioutput_class" if "classifier" in entry["flow_estimator"] or "chain" in entry["flow_estimator"] else "multioutput"
if "n_classes" in names or types.get("y") == "ptr<i32>":
return "classification"
return "regression"


def timed(fn, repeats: int) -> float:
best = float("inf")
for _ in range(repeats):
t0 = time.perf_counter()
fn()
best = min(best, (time.perf_counter() - t0) * 1000.0)
return best


def main() -> int:
ap = argparse.ArgumentParser()
ap.add_argument("--repeats", type=int, default=5)
ap.add_argument("--output", type=Path, default=ROOT / "benchmarks" / "estimators_sklearn.json")
args = ap.parse_args()

registry = json.loads(REGISTRY.read_text())
runnable = [e for e in registry["entries"] if e["bucket"] == "runnable"]
classes = dict(all_estimators())
constructors = _constructors()

iris = load_iris()
diabetes = load_diabetes()
data = {
"classification": (iris.data.astype(np.float64), iris.target.astype(np.float64)),
"unsupervised": (iris.data.astype(np.float64), None),
"regression": (diabetes.data.astype(np.float64), diabetes.target.astype(np.float64)),
}
y_multi = np.column_stack([diabetes.target, diabetes.target * 0.5])
data["multioutput"] = (diabetes.data.astype(np.float64), y_multi)
y_labels = np.column_stack([iris.target, (iris.target + 1) % 3])
data["multioutput_class"] = (iris.data.astype(np.float64), y_labels)

rows = []
for entry in runnable:
name = entry["sklearn_estimator"]
cls = classes.get(name)
if cls is None:
rows.append({"flow_estimator": entry["flow_estimator"], "sklearn_estimator": name,
"status": "unavailable", "reason": "not in sklearn.utils.all_estimators()"})
continue
kind = dataset_kind(entry)
X, y = data[kind]
try:
with warnings.catch_warnings():
warnings.simplefilter("ignore")
build = constructors.get(name)
model = build(cls) if build else cls()
fit_ms = timed((lambda: model.fit(X)) if y is None else (lambda: model.fit(X, y)), args.repeats)
pred_ms = 0.0
for method in ("predict", "transform"):
if hasattr(model, method):
try:
pred_ms = timed(lambda m=method: getattr(model, m)(X), args.repeats)
except Exception:
pred_ms = 0.0
break
rows.append({
"flow_estimator": entry["flow_estimator"],
"sklearn_estimator": name,
"dataset": kind,
"fit_ms": fit_ms,
"pred_ms": pred_ms,
"timing_unit": "ms",
"status": "ok",
})
except Exception as exc: # the estimator refused this data or these defaults
rows.append({
"flow_estimator": entry["flow_estimator"],
"sklearn_estimator": name,
"dataset": kind,
"status": "failed",
"reason": f"{type(exc).__name__}: {exc}"[:200],
})

ok = sum(1 for r in rows if r["status"] == "ok")
args.output.write_text(json.dumps({"schema_version": 1, "repeats": args.repeats,
"counts": {"rows": len(rows), "ok": ok},
"rows": rows}, indent=2) + "\n")
print(f"timed {ok} of {len(rows)} scikit-learn estimators -> {args.output.name}")
for r in rows:
if r["status"] != "ok":
print(f" {r['status']}: {r['flow_estimator']} -> {r['sklearn_estimator']}: {r.get('reason','')[:110]}")
return 0


if __name__ == "__main__":
raise SystemExit(main())
151 changes: 151 additions & 0 deletions benchmarks/compare_estimators.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,151 @@
#!/usr/bin/env python3
"""Join the wide Flow and scikit-learn estimator timings into one artifact.

This is a breadth measurement, and it says so in its own output. The canonical
19 rows carry a parity contract, declared tolerances and a disparity report;
these rows carry none of that. An estimator here is timed on the same data with
each library's own defaults, which answers "is the Flow implementation in the
same performance league" and does not answer "does it compute the same thing".

`speedup` is `sklearn_ms / flow_ms` over fit plus predict, as everywhere else.
"""
from __future__ import annotations

import argparse
import json
import re
from pathlib import Path

ROOT = Path(__file__).resolve().parents[1]
LINE = re.compile(r"^ESTIMATOR\|([a-z_0-9]+)\|([0-9.eE+-]+)\|([0-9.eE+-]+)\|(\d+)\|ok$")

# flow_now_ns resolves to about a microsecond through this harness, so a total
# at or below this is a rounding artifact rather than a measurement.
RESOLUTION_MS = 0.00002


def parse_flow(path: Path) -> dict[str, dict]:
rows: dict[str, dict] = {}
for line in path.read_text().splitlines():
m = LINE.match(line.strip())
if not m:
continue
name, fit, pred, reps = m.group(1), float(m.group(2)), float(m.group(3)), int(m.group(4))
prior = rows.get(name)
# Repeated runs of the same file: keep the fastest, as every other
# harness here does.
if prior is None or fit + pred < prior["fit_ms"] + prior["pred_ms"]:
rows[name] = {"fit_ms": fit, "pred_ms": pred, "repeats": reps}
return rows


def main() -> int:
ap = argparse.ArgumentParser()
ap.add_argument("flow", type=Path, help="captured ESTIMATOR| lines from the generated benchmarks")
ap.add_argument("--sklearn", type=Path, default=ROOT / "benchmarks" / "estimators_sklearn.json")
ap.add_argument("--registry", type=Path, default=ROOT / "benchmarks" / "estimator_coverage.json")
ap.add_argument("--output", type=Path, default=ROOT / "benchmarks" / "estimator_comparison.json")
ap.add_argument("--note", default="", help="recorded verbatim in the artifact, for measurement conditions")
args = ap.parse_args()

flow = parse_flow(args.flow)
sk = {r["flow_estimator"]: r for r in json.loads(args.sklearn.read_text())["rows"]}
registry = json.loads(args.registry.read_text())

rows = []
for entry in registry["entries"]:
name = entry["flow_estimator"]
if entry["bucket"] != "runnable":
row = {"flow_estimator": name, "status": entry["bucket"], "reason": entry.get("reason", "")}
if entry.get("sklearn_estimator"):
row["sklearn_estimator"] = entry["sklearn_estimator"]
# A simplified implementation still gets its times shown, because
# hiding them reads as concealment. They carry no ratio, because a
# ratio against a different algorithm is not a speedup.
f, k = flow.get(name), sk.get(name)
if entry["bucket"] == "simplified" and f and k and k.get("status") == "ok":
row["flow_ms"] = f["fit_ms"] + f["pred_ms"]
row["sklearn_ms"] = k["fit_ms"] + k["pred_ms"]
row["timing_unit"] = "ms"
rows.append(row)
continue
f, s = flow.get(name), sk.get(name)
if f is None:
rows.append({"flow_estimator": name, "sklearn_estimator": entry["sklearn_estimator"],
"status": "flow_missing", "reason": "no ESTIMATOR line in the Flow output"})
continue
if s is None or s.get("status") != "ok":
rows.append({"flow_estimator": name, "sklearn_estimator": entry["sklearn_estimator"],
"status": "sklearn_missing", "reason": (s or {}).get("reason", "absent")})
continue
flow_ms = f["fit_ms"] + f["pred_ms"]
sk_ms = s["fit_ms"] + s["pred_ms"]
# The harness repeats a fast estimator until the clock can see it, so
# this should not trigger any more. It stays as the backstop it always
# was: a Flow side at the floor has been rounded rather than measured.
if flow_ms <= RESOLUTION_MS:
rows.append({
"flow_estimator": name,
"sklearn_estimator": entry["sklearn_estimator"],
"dataset": s["dataset"],
"flow_ms": flow_ms,
"sklearn_ms": sk_ms,
"speedup": None,
"timing_unit": "ms",
"status": "below_resolution",
"reason": f"Flow side at or under {RESOLUTION_MS} ms, which is the clock's floor here",
})
continue
rows.append({
"flow_estimator": name,
"sklearn_estimator": entry["sklearn_estimator"],
"dataset": s["dataset"],
"flow_ms": flow_ms,
"sklearn_ms": sk_ms,
"speedup": (sk_ms / flow_ms) if flow_ms > 0 else None,
"flow_repeats": f["repeats"],
"timing_unit": "ms",
"status": "ok",
})

compared = [r for r in rows if r["status"] == "ok" and r["speedup"] is not None]
unmeasured = [r for r in rows if r["status"] == "below_resolution"]
wins = sum(1 for r in compared if r["speedup"] >= 1.0)
payload = {
"schema_version": 1,
"contract": "breadth timing only; no parity contract, no declared tolerances, "
"each library on its own defaults over the same data",
"measurement": "fastest of several rounds per estimator on both sides, which is "
"what survives a machine that is not idle",
"note": args.note,
"counts": {
"registry_estimators": len(registry["entries"]),
"compared": len(compared),
"flow_wins": wins,
"sklearn_wins": len(compared) - wins,
"below_resolution": len(unmeasured),
"simplified": sum(1 for r in rows if r["status"] == "simplified"),
},
"rows": rows,
}
args.output.write_text(json.dumps(payload, indent=2) + "\n")
print(f"compared {len(compared)} estimators: {wins} Flow wins, {len(compared) - wins} scikit-learn wins")
simplified = [r for r in rows if r["status"] == "simplified"]
if simplified:
print(f"{len(simplified)} rows excluded, implementation declared simplified in its own comments:")
for r in sorted(simplified, key=lambda r: r["flow_estimator"]):
times = (f" flow={r['flow_ms']:.3f} sklearn={r['sklearn_ms']:.3f}"
if "flow_ms" in r else "")
print(f" {r['flow_estimator']:32s}{times}")
if unmeasured:
print(f"{len(unmeasured)} rows left unranked, Flow side under the clock's resolution:")
for r in sorted(unmeasured, key=lambda r: r["flow_estimator"]):
print(f" {r['flow_estimator']:32s} flow={r['flow_ms']:.4f} sklearn={r['sklearn_ms']:9.3f}")
losers = sorted((r for r in compared if r["speedup"] < 1.0), key=lambda r: r["speedup"])
for r in losers:
print(f" {r['speedup']:6.2f}x {r['flow_estimator']:32s} flow={r['flow_ms']:9.3f} sklearn={r['sklearn_ms']:9.3f}")
return 0


if __name__ == "__main__":
raise SystemExit(main())
Loading
Loading