Skip to content
Merged
Show file tree
Hide file tree
Changes from 1 commit
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
1 change: 1 addition & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Comment thread
cursor[bot] marked this conversation as resolved.
Outdated
- `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.
Expand Down
25 changes: 23 additions & 2 deletions src/pixelator/pna/pixeldataset/edgelist.py
Original file line number Diff line number Diff line change
Expand Up @@ -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:
Expand Down
22 changes: 22 additions & 0 deletions tests/pna/pixeldataset/test_anndata_helper.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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))
Loading