From 7cf60ff8a798eafa1dab7d42403e718464d7f168 Mon Sep 17 00:00:00 2001 From: Cursor Agent Date: Mon, 5 Oct 2026 12:12:15 +0000 Subject: [PATCH 1/4] fix(edgelist): rename legacy marker columns in record batches Edgelist.to_record_batches left marker1 and marker2 unchanged, unlike to_polars, to_df, and iterator. Co-authored-by: Adrien Coulier --- CHANGELOG.md | 1 + src/pixelator/pna/pixeldataset/edgelist.py | 25 +++++++++++++++++-- tests/pna/pixeldataset/test_anndata_helper.py | 22 ++++++++++++++++ 3 files changed, 46 insertions(+), 2 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index c2e7f2232..0a0e5cd34 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -54,6 +54,7 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 passed directly to native community detection, and the recovered components are unchanged. ### Fixed +- `Edgelist.to_record_batches` renames legacy `marker1` and `marker2` columns to `marker_1` and `marker_2`, matching `to_polars`, `to_df`, and `iterator`. The stream on dev left those column names unchanged. - `coarsened_pmds_layout` sizes PMDS pivots from the full graph when Leiden yields too few communities, so a valid low `pivots` no longer fails `pmds_layout`'s `0.2 * n` lower bound. diff --git a/src/pixelator/pna/pixeldataset/edgelist.py b/src/pixelator/pna/pixeldataset/edgelist.py index 3bb7e0dab..ced4a8424 100644 --- a/src/pixelator/pna/pixeldataset/edgelist.py +++ b/src/pixelator/pna/pixeldataset/edgelist.py @@ -99,12 +99,33 @@ def to_polars(self) -> pl.DataFrame: def to_record_batches( self, batch_size: int = 1_000_000 ) -> Iterable[pa.RecordBatch]: - """Get the edgelist as a stream of pyarrow RecordBatches.""" + """Get the edgelist as a stream of pyarrow RecordBatches. + + Legacy ``marker1`` and ``marker2`` columns are renamed to + ``marker_1`` and ``marker_2``, matching :meth:`to_polars`. + """ query = self._query_builder.edgelist_query( normalize_input_to_list(self.components) ) with self._view.open() as session: - yield from session.execute_arrow_reader(query=query, batch_size=batch_size) + for batch in session.execute_arrow_reader( + query=query, batch_size=batch_size + ): + names = set(batch.schema.names) + if "marker1" not in names and "marker2" not in names: + yield batch + continue + frame = pl.from_arrow(batch) + if not isinstance(frame, pl.DataFrame): + yield batch + continue + renamed = self._handle_backwards_compatibility(frame.lazy()).collect() + table = renamed.to_arrow() + batches = table.to_batches(max_chunksize=max(batch.num_rows, 1)) + if batches: + yield from batches + else: + yield pa.RecordBatch.from_pylist([], schema=table.schema) def _iterator(self) -> Iterable[tuple[str, pl.LazyFrame]]: with self._view.open() as session: diff --git a/tests/pna/pixeldataset/test_anndata_helper.py b/tests/pna/pixeldataset/test_anndata_helper.py index e5daf6056..5b4b399a6 100644 --- a/tests/pna/pixeldataset/test_anndata_helper.py +++ b/tests/pna/pixeldataset/test_anndata_helper.py @@ -3,8 +3,10 @@ Copyright © 2025 Pixelgen Technologies AB. """ +import shutil from pathlib import Path +import duckdb import numpy as np import polars as pl import pytest @@ -331,3 +333,23 @@ def test_skips_bump_when_prerequisites_are_not_met( assert (adata_old[:, "MarkerC"].X == not_bumped[0][:, "MarkerC"].X).all() assert (adata_new[:, "MarkerC"].X == not_bumped[1][:, "MarkerC"].X).all() + + +def test_record_batches_rename_legacy_marker_columns(tmp_path: Path, pxl_file: Path): + """A legacy edgelist stream uses the same column names as ``to_polars``.""" + legacy = tmp_path / "legacy.pxl" + shutil.copy(pxl_file, legacy) + with duckdb.connect(str(legacy)) as connection: + connection.execute("ALTER TABLE edgelist RENAME COLUMN marker_1 TO marker1") + connection.execute("ALTER TABLE edgelist RENAME COLUMN marker_2 TO marker2") + dataset = PNAPixelDataset.from_files(legacy) + streamed = pl.concat( + [pl.from_arrow(batch) for batch in dataset.edgelist().to_record_batches()], + how="vertical", + ) + loaded = dataset.edgelist().to_polars() + assert "marker_1" in streamed.columns + assert "marker_2" in streamed.columns + assert "marker1" not in streamed.columns + assert "marker2" not in streamed.columns + assert streamed.sort(streamed.columns).equals(loaded.sort(loaded.columns)) From 69c95c0ea4e9f22164b5b58e7a7cd0ade1021484 Mon Sep 17 00:00:00 2001 From: Erik Date: Mon, 5 Oct 2026 17:23:17 +0200 Subject: [PATCH 2/4] Update CHANGELOG.md --- CHANGELOG.md | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 0a0e5cd34..c03a4a9cf 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -54,7 +54,7 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 passed directly to native community detection, and the recovered components are unchanged. ### Fixed -- `Edgelist.to_record_batches` renames legacy `marker1` and `marker2` columns to `marker_1` and `marker_2`, matching `to_polars`, `to_df`, and `iterator`. The stream on dev left those column names unchanged. +- `Edgelist.to_record_batches` renames legacy `marker1` and `marker2` columns to `marker_1` and `marker_2`, matching `to_polars`, `to_df`, and `iterator`. - `coarsened_pmds_layout` sizes PMDS pivots from the full graph when Leiden yields too few communities, so a valid low `pivots` no longer fails `pmds_layout`'s `0.2 * n` lower bound. From 9027084eb42436c4cc0227ca5d5b5ae9eac3c08e Mon Sep 17 00:00:00 2001 From: Adrien Coulier Date: Tue, 6 Oct 2026 10:45:49 +0200 Subject: [PATCH 3/4] Rename columns upstream with SQL query --- src/pixelator/pna/pixeldataset/edgelist.py | 82 ++++++++++++---------- 1 file changed, 45 insertions(+), 37 deletions(-) diff --git a/src/pixelator/pna/pixeldataset/edgelist.py b/src/pixelator/pna/pixeldataset/edgelist.py index ced4a8424..9898edc14 100644 --- a/src/pixelator/pna/pixeldataset/edgelist.py +++ b/src/pixelator/pna/pixeldataset/edgelist.py @@ -11,11 +11,18 @@ import polars as pl import pyarrow as pa -from pixelator.pna.pixeldataset.io import PixelDataViewer, QueryBuilder +from pixelator.pna.pixeldataset.io import ( + PixelDataViewer, + PixelDataViewerSession, + Query, + QueryBuilder, +) from pixelator.pna.pixeldataset.io.anndata_helper import AnnDataHelper from pixelator.pna.pixeldataset.types import Component from pixelator.pna.utils import normalize_input_to_list, normalize_input_to_set +_LEGACY_MARKER_COLUMNS = (("marker1", "marker_1"), ("marker2", "marker_2")) + class Edgelist: """Representation of an edgelist. @@ -56,9 +63,31 @@ def components(self) -> set[str]: ) return set(adata.obs.index.to_list()) - def _handle_backwards_compatibility(self, df: pl.LazyFrame) -> pl.LazyFrame: - # Handle legacy marker names - return df.rename({"marker1": "marker_1", "marker2": "marker_2"}, strict=False) + def _edgelist_query( + self, session: PixelDataViewerSession, components: list[str] | None + ) -> Query: + """Build the edgelist query, aliasing legacy marker column names.""" + query = self._query_builder.edgelist_query(components) + columns = set( + session.execute_eager( + Query( + "SELECT column_name FROM (DESCRIBE SELECT * FROM edgelist)", + {}, + ) + )["column_name"].to_list() + ) + replacements = [ + f"{old} AS {new}" for old, new in _LEGACY_MARKER_COLUMNS if old in columns + ] + if not replacements: + return query + return Query( + sql=( + f"SELECT * RENAME ({', '.join(replacements)}) " + f"FROM ({query.sql}) AS edgelist" + ), + params=query.params, + ) def __len__(self) -> int: """Get the number of edges in the edgelist.""" @@ -74,12 +103,10 @@ def is_empty(self) -> bool: def to_df(self) -> pd.DataFrame: """Get the edgelist as a pandas DataFrame.""" - query = self._query_builder.edgelist_query( - normalize_input_to_list(self.components) - ) + components = normalize_input_to_list(self.components) with self._view.open() as session: df = ( - self._handle_backwards_compatibility(session.execute_lazy(query)) + session.execute_lazy(self._edgelist_query(session, components)) .collect() .to_pandas() ) @@ -87,12 +114,10 @@ def to_df(self) -> pd.DataFrame: def to_polars(self) -> pl.DataFrame: """Get the edgelist as a polars DataFrame.""" - query = self._query_builder.edgelist_query( - normalize_input_to_list(self.components) - ) + components = normalize_input_to_list(self.components) with self._view.open() as session: - df = self._handle_backwards_compatibility( - session.execute_lazy(query) + df = session.execute_lazy( + self._edgelist_query(session, components) ).collect() return df @@ -104,36 +129,19 @@ def to_record_batches( Legacy ``marker1`` and ``marker2`` columns are renamed to ``marker_1`` and ``marker_2``, matching :meth:`to_polars`. """ - query = self._query_builder.edgelist_query( - normalize_input_to_list(self.components) - ) + components = normalize_input_to_list(self.components) with self._view.open() as session: - for batch in session.execute_arrow_reader( - query=query, batch_size=batch_size - ): - names = set(batch.schema.names) - if "marker1" not in names and "marker2" not in names: - yield batch - continue - frame = pl.from_arrow(batch) - if not isinstance(frame, pl.DataFrame): - yield batch - continue - renamed = self._handle_backwards_compatibility(frame.lazy()).collect() - table = renamed.to_arrow() - batches = table.to_batches(max_chunksize=max(batch.num_rows, 1)) - if batches: - yield from batches - else: - yield pa.RecordBatch.from_pylist([], schema=table.schema) + yield from session.execute_arrow_reader( + query=self._edgelist_query(session, components), + batch_size=batch_size, + ) def _iterator(self) -> Iterable[tuple[str, pl.LazyFrame]]: with self._view.open() as session: for component in self.components: - query = self._query_builder.edgelist_query([component]) yield ( component, - session.execute_lazy(query), + session.execute_lazy(self._edgelist_query(session, [component])), ) def iterator(self) -> Iterable[Component]: @@ -148,7 +156,7 @@ def iterator(self) -> Iterable[Component]: # here is that otherwise the object is not pickable, and thus not handled # well by the analysis manager. We should revisit this in the future. component_id=name, - frame=self._handle_backwards_compatibility(df).collect().lazy(), + frame=df.collect().lazy(), ) def __str__(self) -> str: From 0743850e6daafbfd3b77ed2ebfb5069989b079e2 Mon Sep 17 00:00:00 2001 From: Adrien Coulier Date: Tue, 6 Oct 2026 11:07:41 +0200 Subject: [PATCH 4/4] Isolate legacy changes to a separate file --- src/pixelator/pna/pixeldataset/edgelist.py | 48 +++++-------------- .../pna/pixeldataset/legacy/__init__.py | 8 ++++ .../pna/pixeldataset/legacy/marker_columns.py | 41 ++++++++++++++++ 3 files changed, 62 insertions(+), 35 deletions(-) create mode 100644 src/pixelator/pna/pixeldataset/legacy/__init__.py create mode 100644 src/pixelator/pna/pixeldataset/legacy/marker_columns.py diff --git a/src/pixelator/pna/pixeldataset/edgelist.py b/src/pixelator/pna/pixeldataset/edgelist.py index 9898edc14..11a829d7f 100644 --- a/src/pixelator/pna/pixeldataset/edgelist.py +++ b/src/pixelator/pna/pixeldataset/edgelist.py @@ -11,18 +11,12 @@ import polars as pl import pyarrow as pa -from pixelator.pna.pixeldataset.io import ( - PixelDataViewer, - PixelDataViewerSession, - Query, - QueryBuilder, -) +from pixelator.pna.pixeldataset.io import PixelDataViewer, Query, QueryBuilder from pixelator.pna.pixeldataset.io.anndata_helper import AnnDataHelper +from pixelator.pna.pixeldataset.legacy import LegacyMarkerQueryAdapter from pixelator.pna.pixeldataset.types import Component from pixelator.pna.utils import normalize_input_to_list, normalize_input_to_set -_LEGACY_MARKER_COLUMNS = (("marker1", "marker_1"), ("marker2", "marker_2")) - class Edgelist: """Representation of an edgelist. @@ -64,30 +58,10 @@ def components(self) -> set[str]: return set(adata.obs.index.to_list()) def _edgelist_query( - self, session: PixelDataViewerSession, components: list[str] | None + self, adapter: LegacyMarkerQueryAdapter, components: list[str] | None ) -> Query: - """Build the edgelist query, aliasing legacy marker column names.""" - query = self._query_builder.edgelist_query(components) - columns = set( - session.execute_eager( - Query( - "SELECT column_name FROM (DESCRIBE SELECT * FROM edgelist)", - {}, - ) - )["column_name"].to_list() - ) - replacements = [ - f"{old} AS {new}" for old, new in _LEGACY_MARKER_COLUMNS if old in columns - ] - if not replacements: - return query - return Query( - sql=( - f"SELECT * RENAME ({', '.join(replacements)}) " - f"FROM ({query.sql}) AS edgelist" - ), - params=query.params, - ) + """Build the edgelist query, adapted for a legacy marker schema when needed.""" + return adapter.adapt(self._query_builder.edgelist_query(components)) def __len__(self) -> int: """Get the number of edges in the edgelist.""" @@ -105,8 +79,9 @@ def to_df(self) -> pd.DataFrame: """Get the edgelist as a pandas DataFrame.""" components = normalize_input_to_list(self.components) with self._view.open() as session: + adapter = LegacyMarkerQueryAdapter(session) df = ( - session.execute_lazy(self._edgelist_query(session, components)) + session.execute_lazy(self._edgelist_query(adapter, components)) .collect() .to_pandas() ) @@ -116,8 +91,9 @@ def to_polars(self) -> pl.DataFrame: """Get the edgelist as a polars DataFrame.""" components = normalize_input_to_list(self.components) with self._view.open() as session: + adapter = LegacyMarkerQueryAdapter(session) df = session.execute_lazy( - self._edgelist_query(session, components) + self._edgelist_query(adapter, components) ).collect() return df @@ -131,17 +107,19 @@ def to_record_batches( """ components = normalize_input_to_list(self.components) with self._view.open() as session: + adapter = LegacyMarkerQueryAdapter(session) yield from session.execute_arrow_reader( - query=self._edgelist_query(session, components), + query=self._edgelist_query(adapter, components), batch_size=batch_size, ) def _iterator(self) -> Iterable[tuple[str, pl.LazyFrame]]: with self._view.open() as session: + adapter = LegacyMarkerQueryAdapter(session) for component in self.components: yield ( component, - session.execute_lazy(self._edgelist_query(session, [component])), + session.execute_lazy(self._edgelist_query(adapter, [component])), ) def iterator(self) -> Iterable[Component]: diff --git a/src/pixelator/pna/pixeldataset/legacy/__init__.py b/src/pixelator/pna/pixeldataset/legacy/__init__.py new file mode 100644 index 000000000..9989b31c3 --- /dev/null +++ b/src/pixelator/pna/pixeldataset/legacy/__init__.py @@ -0,0 +1,8 @@ +"""Patches that make older ``.pxl`` files look like the current schema. + +Copyright © 2025 Pixelgen Technologies AB. +""" + +from pixelator.pna.pixeldataset.legacy.marker_columns import LegacyMarkerQueryAdapter + +__all__ = ["LegacyMarkerQueryAdapter"] diff --git a/src/pixelator/pna/pixeldataset/legacy/marker_columns.py b/src/pixelator/pna/pixeldataset/legacy/marker_columns.py new file mode 100644 index 000000000..b4cf61e22 --- /dev/null +++ b/src/pixelator/pna/pixeldataset/legacy/marker_columns.py @@ -0,0 +1,41 @@ +"""Adapt queries so an edgelist with legacy marker column names matches the current schema. + +Copyright © 2025 Pixelgen Technologies AB. +""" + +from __future__ import annotations + +from pixelator.pna.pixeldataset.io import PixelDataViewerSession, Query + +_LEGACY_MARKER_COLUMNS = (("marker1", "marker_1"), ("marker2", "marker_2")) + + +class LegacyMarkerQueryAdapter: + """Rename ``marker1`` and ``marker2`` on queries against a pre-rename edgelist. + + A current file is left unchanged. Construct one adapter per open session and + reuse it for every query in that session. + """ + + def __init__(self, session: PixelDataViewerSession) -> None: + """Record which legacy marker columns this session's edgelist still has.""" + columns = set( + session.execute_eager( + Query( + "SELECT column_name FROM (DESCRIBE SELECT * FROM edgelist)", + {}, + ) + )["column_name"].to_list() + ) + self._rename = ", ".join( + f"{old} AS {new}" for old, new in _LEGACY_MARKER_COLUMNS if old in columns + ) + + def adapt(self, query: Query) -> Query: + """Return ``query`` with legacy marker columns renamed, when any are present.""" + if not self._rename: + return query + return Query( + sql=(f"SELECT * RENAME ({self._rename}) FROM ({query.sql}) AS edgelist"), + params=query.params, + )