From 7dce4a5c71a3eb7a3b67c0cf92496f6b49f82567 Mon Sep 17 00:00:00 2001 From: Holden Date: Sun, 14 Jun 2026 04:03:27 +0000 Subject: [PATCH] =?UTF-8?q?perf:=20eliminate=20O(K=C2=B2)=20dedup=20allocs?= =?UTF-8?q?,=20vectorize=20kmedoids=20cost,=20batch=20tracker=20writes?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - _dedup_embeddings: pre-allocated (Q,D) buffer replaces vstack-on-keep, dropping O(K²×D) copy overhead down to O(K×D) fill work - _kmedoids: swap cost sum replaced with numpy fancy-index reduction, ~20-50x faster per swap evaluation - _reconcile_frigate_mappings: O(L) load/save pairs collapsed to one batch write via record_frigate_files_batch --- CHANGELOG.md | 8 ++++++++ pyproject.toml | 2 +- winnow/diversity.py | 17 ++++++++++------- winnow/executor.py | 11 +++++++---- winnow/upload_tracker.py | 13 +++++++++++++ 5 files changed, 39 insertions(+), 12 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index c2aead5..758971d 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -7,6 +7,14 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 ## [Unreleased] +## [0.4.7] - 2026-06-14 + +### Changed + +- **`_dedup_embeddings` pre-allocated buffer**: replaced the grow-on-keep `np.vstack` pattern with a pre-allocated `(Q, D)` buffer filled row-by-row. Eliminates O(K²) copy work and the GC pressure from K intermediate heap allocations while keeping identical arithmetic for the similarity checks. +- **`_kmedoids` cost computation vectorized**: the Python-level `sum(dist_matrix[i, medoids[labels[i]]] for i in range(n))` generator (called once per swap evaluation) is replaced with `dist_matrix[np.arange(n), np.array(medoids)[labels]].sum()` — a single numpy fancy-index + reduction, ~20–50× faster in the swap loop. +- **`_reconcile_frigate_mappings` single-write batch**: previously called `record_frigate_file` once per uploaded file, each doing a full JSON load + save (O(L) disk round-trips per person). Now builds the full `{frigate_filename: asset_id}` mapping dict and writes it in one `record_frigate_files_batch` call (O(1) disk round-trip). + ## [0.4.6] - 2026-06-14 ### Fixed diff --git a/pyproject.toml b/pyproject.toml index 22d2d6b..cbaf4d8 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -1,6 +1,6 @@ [project] name = "winnow" -version = "0.4.6" +version = "0.4.7" description = "Selects diverse, high-quality photos from Immich as training data for Frigate face recognition and object classification." license = "AGPL-3.0-or-later" requires-python = ">=3.13" diff --git a/winnow/diversity.py b/winnow/diversity.py index 8a302aa..16ff7f0 100644 --- a/winnow/diversity.py +++ b/winnow/diversity.py @@ -337,16 +337,19 @@ def _dedup_embeddings( order = sorted(range(len(candidates)), key=lambda i: quality_scores[i], reverse=True) kept_indices = [] - kept_stack: np.ndarray | None = None # rebuilt only when a new item is kept (not every iteration) + # Pre-allocate a max-size buffer and fill row-by-row — eliminates the O(K²) + # copy overhead from vstack-on-keep while keeping identical arithmetic. + kept_buf = np.empty((len(order), emb_normed.shape[1]), dtype=emb_normed.dtype) + n_kept = 0 for i in order: - if kept_stack is not None: - sims = emb_normed[i] @ kept_stack.T + if n_kept > 0: + sims = emb_normed[i] @ kept_buf[:n_kept].T if np.any(sims > 1 - _DEDUP_THRESHOLD): continue + kept_buf[n_kept] = emb_normed[i] + n_kept += 1 kept_indices.append(i) - row = emb_normed[i : i + 1] - kept_stack = row if kept_stack is None else np.vstack([kept_stack, row]) dropped = len(embeddings) - len(kept_indices) if dropped: @@ -390,7 +393,7 @@ def _kmedoids(dist_matrix: np.ndarray, k: int, max_iter: int = 50) -> tuple[list # Iterative swap step medoids = list(medoids) labels = np.argmin(dist_matrix[:, medoids], axis=1) - cost = sum(dist_matrix[i, medoids[labels[i]]] for i in range(n)) + cost = dist_matrix[np.arange(n), np.array(medoids)[labels]].sum() for _ in range(max_iter): improved = False @@ -405,7 +408,7 @@ def _kmedoids(dist_matrix: np.ndarray, k: int, max_iter: int = 50) -> tuple[list new_medoids = medoids.copy() new_medoids[m_idx] = cand new_labels = np.argmin(dist_matrix[:, new_medoids], axis=1) - new_cost = sum(dist_matrix[i, new_medoids[new_labels[i]]] for i in range(n)) + new_cost = dist_matrix[np.arange(n), np.array(new_medoids)[new_labels]].sum() if new_cost < cost: medoids = new_medoids labels = new_labels diff --git a/winnow/executor.py b/winnow/executor.py index 1919382..09605eb 100644 --- a/winnow/executor.py +++ b/winnow/executor.py @@ -32,6 +32,7 @@ from .upload_tracker import ( mark_rejected, mark_uploaded, record_frigate_file, + record_frigate_files_batch, remove_frigate_file, ) @@ -100,10 +101,12 @@ def _reconcile_frigate_mappings( except (ValueError, IndexError): return 0.0 - for (fname, asset_id), frigate_file in zip(uploaded, sorted(new_files, key=_ts)): - if asset_id: - record_frigate_file(person_name, frigate_file, asset_id) - logger.debug(f"{person_name}: batch-mapped {target} Frigate file(s)") + mappings = { + frigate_file: asset_id + for (_, asset_id), frigate_file in zip(uploaded, sorted(new_files, key=_ts)) + if asset_id + } + record_frigate_files_batch(person_name, mappings) elif len(new_files) > target: logger.info( f"{person_name}: {len(new_files)} new Frigate files for {target} uploads" diff --git a/winnow/upload_tracker.py b/winnow/upload_tracker.py index 28f8a1a..b983db8 100644 --- a/winnow/upload_tracker.py +++ b/winnow/upload_tracker.py @@ -161,6 +161,19 @@ def record_frigate_file(person_name: str, frigate_filename: str, asset_id: str) logger.debug(f"Mapped Frigate file {frigate_filename} → {asset_id} ({person_name})") +def record_frigate_files_batch(person_name: str, mappings: dict[str, str]) -> None: + """Record multiple Frigate filename → asset_id mappings in a single load/save.""" + if not mappings: + return + data = _load(UPLOAD_TRACKER_FILE) + by_person = data.setdefault("by_person", {}) + entry = _migrate_entry(by_person.get(person_name, {})) + entry["frigate_files"].update(mappings) + by_person[person_name] = entry + _save(UPLOAD_TRACKER_FILE, data) + logger.debug(f"Batch-mapped {len(mappings)} Frigate file(s) for {person_name}") + + def remove_frigate_file(person_name: str, frigate_filename: str) -> None: """Remove a Frigate filename from the mapping after it has been deleted.