From 660a2878e47088e6cd0539709df54cabb79a914e Mon Sep 17 00:00:00 2001 From: Francesc Alted Date: Thu, 16 Jul 2026 12:05:21 +0200 Subject: [PATCH 1/8] Fix engine=blosc2.jit against pandas 3.0.3 and implement Series.map DataFrame.apply(f, engine=blosc2.jit) returned a raw NumPy array instead of a properly indexed DataFrame/Series when raw=False (the default), since pandas only reconstructs the pandas object itself for raw=True. Series.map(f, engine=blosc2.jit) now works instead of always raising NotImplementedError. Non-numeric columns raise a clear ValueError instead of a deep numexpr error. Co-Authored-By: Claude Sonnet 5 --- RELEASE_NOTES.md | 11 +++++ bench/bench_pandas_engine.py | 73 ++++++++++++++++++++++++++++ doc/guides/index.rst | 1 + doc/guides/pandas_engine.md | 83 ++++++++++++++++++++++++++++++++ plans/enhancing-ctable-phase2.md | 42 ++++++++++++++++ src/blosc2/proxy.py | 47 +++++++++++++----- tests/test_pandas_udf_engine.py | 65 ++++++++++++++++++++++++- 7 files changed, 309 insertions(+), 13 deletions(-) create mode 100644 bench/bench_pandas_engine.py create mode 100644 doc/guides/pandas_engine.md diff --git a/RELEASE_NOTES.md b/RELEASE_NOTES.md index c2b73a4ca..a96713a15 100644 --- a/RELEASE_NOTES.md +++ b/RELEASE_NOTES.md @@ -4,6 +4,17 @@ XXX version-specific blurb XXX +### Bug fixes + +- Fixed `engine=blosc2.jit` for `DataFrame.apply` against pandas 3.0.3: with + the default `raw=False`, the engine returned a raw NumPy array instead of + a properly indexed `DataFrame`/`Series`, so results only matched plain + `apply()` by value, never by type. `Series.map(func, engine=blosc2.jit)` + is now implemented (it previously always raised `NotImplementedError`). + Non-numeric columns now raise a clear `ValueError` instead of a deep + `numexpr` error. See the new guide "Using Blosc2 as a pandas engine" and + `bench/bench_pandas_engine.py`. + ## Changes from 4.8.0 to 4.8.1 ### Improvements diff --git a/bench/bench_pandas_engine.py b/bench/bench_pandas_engine.py new file mode 100644 index 000000000..b73a1b083 --- /dev/null +++ b/bench/bench_pandas_engine.py @@ -0,0 +1,73 @@ +####################################################################### +# Copyright (c) 2019-present, Blosc Development Team +# All rights reserved. +# +# SPDX-License-Identifier: BSD-3-Clause +####################################################################### + +# Benchmark: DataFrame.apply(f, engine=blosc2.jit) vs plain DataFrame.apply(f) +# +# engine=blosc2.jit calls the vectorized function once per column (the +# default axis=0), so the win comes from the Blosc2/numexpr compute engine +# (operator fusion, multi-threading) beating plain NumPy on a +# multi-operation elementwise expression over a large 1D array. This script +# measures that on a 1,000,000-row, 8-column frame. +# +# Note: axis=1 (row-wise) is NOT a good fit for this engine. It still calls +# the function once per row in a Python loop either way, and for a handful +# of columns the wrapping overhead per call (building a compute-engine proxy +# for a tiny array) is larger than the win, so engine=blosc2.jit is actually +# *slower* than plain apply(axis=1) in that case. Use axis=0 (or restructure +# the computation to operate on whole columns) to get the engine's benefit. +# +# Each measurement is the minimum of NRUNS repetitions to reduce noise. + +from time import perf_counter + +import numpy as np +import pandas as pd + +import blosc2 + +NRUNS = 3 +NROWS = 1_000_000 +NCOLS = 8 + + +def make_df(): + rng = np.random.default_rng(0) + return pd.DataFrame( + {f"c{i}": rng.random(NROWS) for i in range(NCOLS)}, + ) + + +def transform(col): + return np.sin(col) * np.cos(col) + col**2 - np.sqrt(np.abs(col)) + np.exp(-col) + + +def timeit(fn): + best = float("inf") + result = None + for _ in range(NRUNS): + t0 = perf_counter() + result = fn() + best = min(best, perf_counter() - t0) + return best, result + + +def main(): + df = make_df() + + t_plain, result_plain = timeit(lambda: df.apply(transform)) + t_engine, result_engine = timeit(lambda: df.apply(transform, engine=blosc2.jit)) + + pd.testing.assert_frame_equal(result_engine, result_plain) + + print(f"rows={NROWS}, cols={NCOLS}") + print(f"plain df.apply(f): {t_plain:.4f} s") + print(f"df.apply(f, engine=blosc2.jit): {t_engine:.4f} s") + print(f"speedup: {t_plain / t_engine:.1f}x") + + +if __name__ == "__main__": + main() diff --git a/doc/guides/index.rst b/doc/guides/index.rst index 7cdfeabef..5022686cd 100644 --- a/doc/guides/index.rst +++ b/doc/guides/index.rst @@ -13,6 +13,7 @@ Topics optimization_tips sharing_across_processes + pandas_engine Command line tools ------------------ diff --git a/doc/guides/pandas_engine.md b/doc/guides/pandas_engine.md new file mode 100644 index 000000000..2b95424b5 --- /dev/null +++ b/doc/guides/pandas_engine.md @@ -0,0 +1,83 @@ +# Using Blosc2 as a pandas Engine + +pandas' `DataFrame.apply` and `Series.map` accept an `engine=` argument: a +callable exposing a `__pandas_udf__` attribute that pandas dispatches to +instead of running the Python-level per-row/per-element loop. `blosc2.jit` +is such an engine. + +The contract is different from a plain `apply`/`map` callback: the function +passed to `engine=blosc2.jit` must be **vectorized** — it is called *once* +with a full NumPy array (a column, a row, or the whole array, depending on +`axis`), not once per element. This is the same contract `@blosc2.jit` +already has everywhere else in the library; using it as a pandas engine +just changes who supplies the array. + +```python +import numpy as np +import pandas as pd +import blosc2 + +df = pd.DataFrame( + { + "a": np.arange(1_000_000, dtype=np.float64), + "b": np.arange(1_000_000, dtype=np.float64), + } +) + + +@blosc2.jit +def add_one(col): + return col + 1 + + +result = df.apply(add_one, engine=blosc2.jit) +``` + +`axis=0` (the default) calls the function once per column; `axis=1` calls it +once per row. **Use `axis=0`** (or restructure the computation so it works +column-wise): the win comes from the Blosc2/numexpr compute engine (operator +fusion, multi-threading) processing one large 1D array per call, and that +only happens for columns. `axis=1` still calls the function once per row — +same as plain pandas — and for a handful of columns, the overhead of +wrapping each tiny row array for the compute engine outweighs any benefit, +so `engine=blosc2.jit` with `axis=1` is typically *slower* than plain +`apply(axis=1)`. See the benchmark below. + +`Series.map(func, engine=blosc2.jit)` works the same way: `func` is called +once with the Series' full underlying array. + +## Limitations + +- Only numeric dtypes are supported. A non-numeric (e.g. object-dtype or + string) column raises a `ValueError` naming the limitation rather than + attempting the computation. +- `na_action="ignore"` is not supported for `map` and raises + `NotImplementedError` — the vectorized-call contract means there is no + per-element step at which to skip a value. +- `Series.apply(func, engine=...)` and `DataFrame.map(func, engine=...)` do + not reach `blosc2.jit` at all: pandas 3's `Series.apply` does not accept + an `engine` keyword for non-string functions, and `DataFrame.map` doesn't + forward `engine` to a dispatch mechanism at all. These are limitations of + the pandas-side API surface, not of the Blosc2 engine. The two entry + points that do reach the engine are `DataFrame.apply` and `Series.map`. + +## Benchmark + +`bench/bench_pandas_engine.py` compares `df.apply(f, engine=blosc2.jit)` +against plain `df.apply(f)` (`axis=0`, the default) on a 1,000,000-row, +8-column frame, for a multi-operation elementwise expression +(`sin(x)*cos(x) + x**2 - sqrt(|x|) + exp(-x)`). Measured on the development +machine (Apple M4, conda env with pandas 3.0.3): + +``` +rows=1000000, cols=8 +plain df.apply(f): 0.1114 s +df.apply(f, engine=blosc2.jit): 0.0260 s +speedup: 4.3x +``` + +Run the script for the numbers on your machine: + +``` +python bench/bench_pandas_engine.py +``` diff --git a/plans/enhancing-ctable-phase2.md b/plans/enhancing-ctable-phase2.md index 381cdeb98..962a96f0a 100644 --- a/plans/enhancing-ctable-phase2.md +++ b/plans/enhancing-ctable-phase2.md @@ -196,6 +196,48 @@ the adapter directly; ADD end-to-end tests that go through real pandas Acceptance: the P1.3 tests green against pandas 3.0.3; no pandas version pin added to any requirements file. +### Implementation notes (landed) + +Empirically the crash described at the top of this section did not +reproduce verbatim against the installed pandas 3.0.3 — `_ensure_numpy_data` +already existed and handled the `.values` fallback. What *did* reproduce, +reading `pandas/core/frame.py`'s `DataFrame.apply` engine dispatch directly: + +- `raw` defaults to `False`, in which case pandas hands the engine the + **DataFrame itself** (not `.values`) and, unlike the `raw=True` path, + does NOT reconstruct a `DataFrame`/`Series` from the result — it returns + whatever the engine gives back verbatim. `PandasUdfEngine.apply` was + returning a raw `ndarray` in this (default!) case, so + `df.apply(f, engine=blosc2.jit)` produced the right values with the wrong + type — silently broken for any code chaining further DataFrame methods + on the result. Fixed by reconstructing the `DataFrame`/`Series` ourselves + (mirroring pandas' own `raw=True` reconstruction code) whenever the input + we received was the original pandas object rather than a raw array. +- `DataFrame.map(func, engine=...)` does not forward `engine` to any + dispatch mechanism at all in pandas 3.0.3 (`DataFrame.map`'s signature + doesn't accept it; it silently becomes a keyword arg forwarded to `func`, + raising `TypeError`). `Series.apply(func, engine=...)` similarly never + reaches `__pandas_udf__` — only `DataFrame.apply` and `Series.map` do. + Documented as pandas-side limitations, not tested as if they were ours. +- `Series.map(engine=blosc2.jit)` now implemented (P1.2): pandas wraps a + raw array result back into a `Series` itself for `map`, so no + reconstruction is needed on our side. +- Fixed a latent bug in `_ensure_numpy_data`'s error message (missing `f` + prefix, wrong attribute — `data.__name__` instead of + `type(data).__name__`) while adding the new numeric-dtype check. +- Benchmark (`bench/bench_pandas_engine.py`, Apple M4): the "point of the + engine" claim in this section's original framing (skip the per-row Python + loop, `axis=1`) does not hold — `axis=1` still calls the function once per + row either way, and for a handful of columns the per-call proxy-wrapping + overhead makes `engine=blosc2.jit` *slower* than plain `apply(axis=1)`. + The real, verified win is on `axis=0` (the default): 4.3x on a + 1,000,000-row/8-column frame with a multi-op elementwise expression + (numexpr operator fusion + threading beating plain NumPy on one large 1D + array per column). Documented this correction in the new guide page. +- Commit: see `git log` for the commit landing this section (message + references P1 by PR title, not inside code/docs, per the cross-cutting + rules). + --- ## P2 — `CTable.assign()` chaining + unbound `blosc2.col()` diff --git a/src/blosc2/proxy.py b/src/blosc2/proxy.py index 0e55eb677..7f346d36d 100644 --- a/src/blosc2/proxy.py +++ b/src/blosc2/proxy.py @@ -864,9 +864,15 @@ def _ensure_numpy_data(data): data = data.values except AttributeError as err: raise ValueError( - "blosc2.jit received an object of type {data.__name__}, which is not supported. " - "Try casting your Series or DataFrame to a NumPy dtype." + f"blosc2.jit received an object of type {type(data).__name__}, which is not " + "supported. Try casting your Series or DataFrame to a NumPy dtype." ) from err + if data.dtype.kind not in "biufc": + raise ValueError( + f"blosc2.jit requires a numeric dtype, got {data.dtype!r}. The Blosc2 engine only " + "supports vectorized numeric computations; cast non-numeric columns before using " + "engine=blosc2.jit." + ) return data @classmethod @@ -874,10 +880,14 @@ def map(cls, data, func, args, kwargs, decorator, skip_na): """ JIT a NumPy array element-wise. In the case of Blosc2, functions are expected to be vectorized NumPy operations, so the function is called - with the NumPy array as the function parameter, instead of calling the - function once for each element. + once with the whole NumPy array, instead of calling the function once + for each element. """ - raise NotImplementedError("The Blosc2 engine does not support map. Use apply instead.") + if skip_na: + raise NotImplementedError("The Blosc2 engine does not support na_action='ignore' in map.") + values = cls._ensure_numpy_data(data) + func = decorator(func) + return func(values, *args, **kwargs) @classmethod def apply(cls, data, func, args, kwargs, decorator, axis): @@ -887,21 +897,34 @@ def apply(cls, data, func, args, kwargs, decorator, axis): with the NumPy array as the function parameter, instead of calling the function once for each column or row. """ - data = cls._ensure_numpy_data(data) + orig = data + values = cls._ensure_numpy_data(data) func = decorator(func) - if data.ndim == 1 or axis is None: + if values.ndim == 1 or axis is None: # pandas Series.apply or pipe - return func(data, *args, **kwargs) + result = func(values, *args, **kwargs) elif axis in (0, "index"): # pandas apply(axis=0) column-wise - result = [func(data[:, row_idx], *args, **kwargs) for row_idx in range(data.shape[1])] - return np.vstack(result).transpose() + result = [func(values[:, col_idx], *args, **kwargs) for col_idx in range(values.shape[1])] + result = np.vstack(result).transpose() elif axis in (1, "columns"): # pandas apply(axis=1) row-wise - result = [func(data[col_idx, :], *args, **kwargs) for col_idx in range(data.shape[0])] - return np.vstack(result) + result = [func(values[row_idx, :], *args, **kwargs) for row_idx in range(values.shape[0])] + result = np.vstack(result) else: raise NotImplementedError(f"Unknown axis '{axis}'. Use one of 0, 1 or None.") + # pandas only reconstructs a DataFrame/Series for us when it called us + # with `raw=True` data (a plain ndarray); when it handed us the + # original DataFrame (`raw=False`, the default), we must return a + # properly indexed pandas object ourselves, mirroring what pandas' + # own raw=True code path does. + if isinstance(result, np.ndarray) and hasattr(orig, "columns"): + if result.ndim == 2: + return orig.__class__(result, index=orig.index, columns=orig.columns) + agg_axis = orig._get_agg_axis(orig._get_axis_number(axis)) + return orig._constructor_sliced(result, index=agg_axis) + return result + jit.__pandas_udf__ = PandasUdfEngine diff --git a/tests/test_pandas_udf_engine.py b/tests/test_pandas_udf_engine.py index e0e322925..cc576c996 100644 --- a/tests/test_pandas_udf_engine.py +++ b/tests/test_pandas_udf_engine.py @@ -18,6 +18,22 @@ def add_one(x): data = np.array([1, 2]) + result = blosc2.jit.__pandas_udf__.map( + data, + add_one, + args=(), + kwargs={}, + decorator=blosc2.jit, + skip_na=False, + ) + assert np.array_equal(result, np.array([2, 3])) + + def test_map_skip_na_not_supported(self): + def add_one(x): + return x + 1 + + data = np.array([1, 2]) + with pytest.raises(NotImplementedError): blosc2.jit.__pandas_udf__.map( data, @@ -25,7 +41,7 @@ def add_one(x): args=(), kwargs={}, decorator=blosc2.jit, - skip_na=False, + skip_na=True, ) def test_apply_1d(self): @@ -117,3 +133,50 @@ def add_one(x): ) expected = np.array([[2, 3, 4], [5, 6, 7]]) assert np.array_equal(result, expected) + + +try: + import pandas as pd + + _pandas_too_old = pd.__version__ < "3" +except ImportError: + pd = None + _pandas_too_old = False + + +@pytest.mark.skipif(pd is None, reason="pandas not installed") +@pytest.mark.skipif(_pandas_too_old, reason="engine= integration targets pandas 3.x") +class TestPandasEngineEndToEnd: + """Exercises engine=blosc2.jit through real pandas, not the adapter directly.""" + + def test_apply_axis0_matches_default_engine(self): + df = pd.DataFrame({"a": [1.0, 2.0, 3.0], "b": [4.0, 5.0, 6.0]}) + expected = df.apply(lambda x: x + 1) + result = df.apply(lambda x: x + 1, engine=blosc2.jit) + pd.testing.assert_frame_equal(result, expected) + + def test_apply_axis1_matches_default_engine(self): + df = pd.DataFrame({"a": [1.0, 2.0, 3.0], "b": [4.0, 5.0, 6.0]}) + expected = df.apply(lambda x: x + 1, axis=1) + result = df.apply(lambda x: x + 1, engine=blosc2.jit, axis=1) + pd.testing.assert_frame_equal(result, expected) + + def test_apply_args_and_kwargs_forwarded(self): + def add_numbers(x, num1, num2=0): + return x + num1 + num2 + + df = pd.DataFrame({"a": [1.0, 2.0], "b": [3.0, 4.0]}) + expected = df.apply(add_numbers, args=(10,), num2=100) + result = df.apply(add_numbers, engine=blosc2.jit, args=(10,), num2=100) + pd.testing.assert_frame_equal(result, expected) + + def test_series_map_matches_default_engine(self): + s = pd.Series([1.0, 2.0, 3.0]) + expected = s.map(lambda x: x + 1) + result = s.map(lambda x: x + 1, engine=blosc2.jit) + pd.testing.assert_series_equal(result, expected) + + def test_apply_object_dtype_raises_clear_error(self): + df = pd.DataFrame({"a": ["x", "y"]}) + with pytest.raises(ValueError, match="numeric dtype"): + df.apply(lambda x: x + 1, engine=blosc2.jit) From 14657fee54c43ccc8ee760d84591c740d344701d Mon Sep 17 00:00:00 2001 From: Francesc Alted Date: Thu, 16 Jul 2026 12:15:52 +0200 Subject: [PATCH 2/8] Add CTable.assign() and unbound blosc2.col() for pandas-3 style chaining col(name) builds a deferred column expression that replays operators against a table's columns only once bound (assign(), indexing, where()), reusing all existing Column/NullableExpr null-propagation and comparison semantics rather than a second expression engine. assign() returns a view with additional computed columns without mutating the table or copying column data. Also fixes a pre-existing bug found while testing the chain end-to-end: CTable.head()/tail() ignored _cached_live_positions and silently returned rows in physical order instead of sorted order after a lazily-sorted view. Co-Authored-By: Claude Sonnet 5 --- RELEASE_NOTES.md | 15 ++ plans/enhancing-ctable-phase2.md | 35 +++++ src/blosc2/__init__.py | 2 + src/blosc2/ctable.py | 240 ++++++++++++++++++++++++++++++- tests/ctable/test_col_expr.py | 132 +++++++++++++++++ 5 files changed, 419 insertions(+), 5 deletions(-) create mode 100644 tests/ctable/test_col_expr.py diff --git a/RELEASE_NOTES.md b/RELEASE_NOTES.md index a96713a15..bb7c95b5a 100644 --- a/RELEASE_NOTES.md +++ b/RELEASE_NOTES.md @@ -4,8 +4,23 @@ XXX version-specific blurb XXX +### New features + +- `CTable.assign(**named_exprs)`: return a view with additional computed + columns, without mutating the table or copying column data. Pairs with + the new `blosc2.col(name)` — an unbound column expression that defers + operator replay until it's bound to a table (`assign()`, `t[...]`, + `where()`) — to write pandas-3-style chains: + `t.assign(profit=col("revenue") - col("cost"))[col("profit") > 0].sort_by("profit", ascending=False).head(10)`. + ### Bug fixes +- Fixed `CTable.head()`/`tail()` silently discarding row order when called + on a lazily-sorted view (e.g. `t.sort_by("col", ascending=False)` on a + view, or any `.sort_by()` result chained off a prior filter): they + ignored `_cached_live_positions` and built a plain physical-order mask + instead, so `t.where(...).sort_by("x", ascending=False).head(10)` came + back in the wrong order. - Fixed `engine=blosc2.jit` for `DataFrame.apply` against pandas 3.0.3: with the default `raw=False`, the engine returned a raw NumPy array instead of a properly indexed `DataFrame`/`Series`, so results only matched plain diff --git a/plans/enhancing-ctable-phase2.md b/plans/enhancing-ctable-phase2.md index 962a96f0a..7f883cdfa 100644 --- a/plans/enhancing-ctable-phase2.md +++ b/plans/enhancing-ctable-phase2.md @@ -415,6 +415,41 @@ Acceptance: all green; `assign` copies no column data (assert via storage identity: the assigned table's `_cols` is the parent's `_cols` object — skip this assertion if the escape valve was taken). +### Implementation notes (landed) + +Landed as designed, no escape valve needed: `assign()` builds one +`CTable._make_view(self, self._valid_rows)` and gives it its own +`_computed_cols`/`col_names`/`_col_widths` copies (the only three structures +`add_computed_column` mutates), then registers each new column directly on +the view's copies — bypassing `add_computed_column`'s own +"cannot add to a view" guard, which is correct for the base table but not +for this internal, view-only registration path. + +`ColExpr` values, plus `Column`/`NullableExpr` values, are bound/unwrapped +against `self` (not the new view) **before** any of the call's new columns +are registered, so a later keyword genuinely cannot see an earlier one in +the same `assign()` call — it fails with the normal unknown-column error, +matching the plan's documented restriction (not accidentally, by construction). + +While testing the Goal section's exact chain end-to-end, found and fixed a +**pre-existing, unrelated bug**: `CTable.head()`/`tail()` build a boolean +mask from `_valid_rows` and ignore `_cached_live_positions`, so calling +`.head(N)` after a lazy `sort_by()` view (`self.base is not None`, always +lazy per that method's own docstring) silently discarded the sort order and +returned rows in physical order instead. Reproduced with plain +`add_computed_column`/`Column` filtering, no `col()`/`assign()` involved — +confirms it predates this item. Fixed by taking `_cached_live_positions[:N]` +/ `[-N:]` through `_view_from_positions` (the same pattern already used by +`_materialize_row` and `_display_positions`) whenever that attribute is set, +before falling back to the existing mask-based fast path. Without this fix, +the plan's own headline example +(`t.assign(...)[...].sort_by(...).head(10)`) silently returned rows in the +wrong order. + +All 11 new tests in `tests/ctable/test_col_expr.py` pass; full +`tests/ctable` (1319 tests) and `tests/ndarray` (4385 tests) suites pass +with no regressions from the `head`/`tail` fix. + --- ## P3 — First-class variable-length string columns (the real project) diff --git a/src/blosc2/__init__.py b/src/blosc2/__init__.py index 23904218a..bb4332d31 100644 --- a/src/blosc2/__init__.py +++ b/src/blosc2/__init__.py @@ -639,6 +639,7 @@ def _raise(exc): NestedColumn, NullPolicy, RowTransformer, + col, ctable_from_cframe, get_null_policy, get_printoptions, @@ -840,6 +841,7 @@ def _raise(exc): "CTableGroupBy", "NestedColumn", "RowTransformer", + "col", "ctable_from_cframe", "Batch", "BatchArray", diff --git a/src/blosc2/ctable.py b/src/blosc2/ctable.py index 17aa2882b..a3da2b840 100644 --- a/src/blosc2/ctable.py +++ b/src/blosc2/ctable.py @@ -16,6 +16,7 @@ import copy import dataclasses import json +import operator import os import pprint import re @@ -3477,6 +3478,111 @@ def finalize(self, expected_size: int) -> np.ndarray | None: return np.concatenate(parts) +class ColExpr: + """Unbound column expression: a recipe that, given a table, evaluates + against that table's columns. + + ``blosc2.col("x") + 1`` builds a deferred computation; passing it to + :meth:`CTable.assign` or using it to index/filter a table binds it, + replaying the operators on ``table["x"]`` — so all :class:`Column` + semantics (null propagation, SQL comparison rules, dictionary/timestamp + handling) apply identically to the bound form ``t.x + 1``. + + Only operators are supported; method calls such as ``col("x").sum()`` + are not, since there is no table to evaluate against yet. Use the bound + form (``t.x.sum()``) for those. + """ + + def __init__(self, bind, repr_str): + self._bind = bind # callable: CTable -> Column/LazyExpr/NullableExpr/scalar + self._repr = repr_str + + def __repr__(self): + return self._repr + + def __getattr__(self, name): + raise AttributeError( + f"{name!r} is not supported on an unbound column expression (only operators are). " + f"Bind it to a table first, e.g. use t.{self._repr}.{name}(...) instead of " + f"{self._repr}.{name}(...)." + ) + + +def _colexpr_binop(op, symbol, reflected=False): + def bind(self, other, t): + left = self._bind(t) + right = other._bind(t) if isinstance(other, ColExpr) else other + return op(right, left) if reflected else op(left, right) + + def method(self, other): + other_repr = other._repr if isinstance(other, ColExpr) else repr(other) + if reflected: + repr_str = f"({other_repr} {symbol} {self._repr})" + else: + repr_str = f"({self._repr} {symbol} {other_repr})" + return ColExpr(lambda t: bind(self, other, t), repr_str) + + return method + + +def _colexpr_unary(op, symbol): + def method(self): + return ColExpr(lambda t: op(self._bind(t)), f"({symbol}{self._repr})") + + return method + + +_COLEXPR_BINOPS = [ + ("__add__", operator.add, "+", False), + ("__radd__", operator.add, "+", True), + ("__sub__", operator.sub, "-", False), + ("__rsub__", operator.sub, "-", True), + ("__mul__", operator.mul, "*", False), + ("__rmul__", operator.mul, "*", True), + ("__truediv__", operator.truediv, "/", False), + ("__rtruediv__", operator.truediv, "/", True), + ("__floordiv__", operator.floordiv, "//", False), + ("__rfloordiv__", operator.floordiv, "//", True), + ("__mod__", operator.mod, "%", False), + ("__rmod__", operator.mod, "%", True), + ("__pow__", operator.pow, "**", False), + ("__rpow__", operator.pow, "**", True), + ("__and__", operator.and_, "&", False), + ("__rand__", operator.and_, "&", True), + ("__or__", operator.or_, "|", False), + ("__ror__", operator.or_, "|", True), + ("__lt__", operator.lt, "<", False), + ("__le__", operator.le, "<=", False), + ("__gt__", operator.gt, ">", False), + ("__ge__", operator.ge, ">=", False), + ("__eq__", operator.eq, "==", False), + ("__ne__", operator.ne, "!=", False), +] +for _dunder, _op, _symbol, _reflected in _COLEXPR_BINOPS: + setattr(ColExpr, _dunder, _colexpr_binop(_op, _symbol, _reflected)) +ColExpr.__neg__ = _colexpr_unary(operator.neg, "-") +ColExpr.__invert__ = _colexpr_unary(operator.invert, "~") +del _dunder, _op, _symbol, _reflected + + +def col(name: str) -> ColExpr: + """Build an unbound column expression referencing a column by name. + + The name is resolved against a table's columns only when the expression + is bound — passed to :meth:`CTable.assign`, or used to index/filter a + table (``t[col("x") > 0]``, ``t.where(col("x") > 0)``). An unknown name + therefore fails at bind time with the table's normal unknown-column + error, not when ``col()`` is called. + + Examples + -------- + >>> import blosc2 + >>> from blosc2 import col + >>> t.assign(profit=col("revenue") - col("cost")) # doctest: +SKIP + """ + return ColExpr(lambda t: t[name], name) + + class CTable(_CTableIndexingMixin, Generic[RowT]): """Columnar compressed table with typed columns and row-oriented access.""" @@ -5847,6 +5953,12 @@ def head(self, N: int = 5) -> CTable: """Return a view of the first *N* live rows (default 5).""" if N <= 0: return self.view(blosc2.zeros(shape=len(self._valid_rows), dtype=np.bool_)) + _slp = getattr(self, "_cached_live_positions", None) + if _slp is not None and self.base is not None: + # A lazily-sorted view: physical row order is not row order, so the + # first N *logical* rows are the first N entries of the stored + # permutation, not the first N physical positions. + return self._view_from_positions(_slp[:N]) if self._n_rows <= N: return self.view(self._valid_rows) @@ -5868,6 +5980,10 @@ def tail(self, N: int = 5) -> CTable: """Return a view of the last *N* live rows (default 5).""" if N <= 0: return self.view(blosc2.zeros(shape=len(self._valid_rows), dtype=np.bool_)) + _slp = getattr(self, "_cached_live_positions", None) + if _slp is not None and self.base is not None: + # See head(): physical order is not row order for a sorted view. + return self._view_from_positions(_slp[-N:] if len(_slp) > N else _slp) if self._n_rows <= N: return self.view(self._valid_rows) @@ -10221,6 +10337,116 @@ def drop_computed_column(self, name: str) -> None: if isinstance(self._storage, FileTableStorage): self._storage.save_schema(self._schema_dict_with_computed()) + @staticmethod + def _coerce_assign_operand(value): + """Reduce an assign() value to a form add_computed_column's transformer + machinery accepts: a LazyExpr, DSLKernel, callable, or string.""" + if isinstance(value, NullableExpr): + return value._expr + if isinstance(value, Column): + value._ensure_queryable() + raw = value._raw_col + return raw if isinstance(raw, blosc2.LazyExpr) else blosc2.lazyexpr(raw) + return value + + def assign(self, **named_exprs) -> CTable: + """Return a view with additional computed columns, without copying data. + + Each keyword argument names a new computed column; the value defines + it, in any of the forms :meth:`add_computed_column` accepts (a string + expression, a :class:`blosc2.LazyExpr`), plus a :class:`Column`, + :class:`NullableExpr`, or an unbound :func:`col` expression. + + Unlike :meth:`add_computed_column`, which mutates the table in place + and cannot be called on a view, ``assign()`` never mutates ``self``: + it returns a new view sharing this table's column storage, with its + own computed-column metadata. This makes it composable in a chain:: + + result = ( + t.assign(profit=col("revenue") - col("cost"))[col("profit") > 0] + .sort_by("profit", ascending=False) + .head(10) + ) + + Values are bound against the *new view* being built, one keyword at a + time, but against a snapshot taken before any of this call's new + columns were added — so a later keyword cannot reference an earlier + one from the same ``assign()`` call (that raises the usual unknown- + column error). Chain two ``assign()`` calls for that:: + + t2 = t.assign(profit=col("revenue") - col("cost")) + t3 = t2.assign(margin=col("profit") / col("revenue")) + + Parameters + ---------- + **named_exprs: + One computed-column definition per keyword. + + Returns + ------- + CTable + A read-only view (see :meth:`select`) with the additional + computed columns. + + Raises + ------ + ValueError + If a name collides with an existing stored or computed column. + + Examples + -------- + >>> import blosc2 + >>> from blosc2 import col + >>> from dataclasses import dataclass + >>> @dataclass + ... class Row: + ... revenue: float = blosc2.field(blosc2.float64()) + ... cost: float = blosc2.field(blosc2.float64()) + >>> t = blosc2.CTable(Row, new_data=[(100.0, 40.0), (50.0, 60.0)]) + >>> t2 = t.assign(profit=col("revenue") - col("cost")) + >>> t2.profit[:] + array([ 60., -10.]) + """ + bound = {} + for name, expr in named_exprs.items(): + _validate_column_name(name) + if name in self._cols: + raise ValueError(f"A stored column named {name!r} already exists.") + if name in self._computed_cols or name in bound: + raise ValueError(f"A computed column named {name!r} already exists.") + value = expr._bind(self) if isinstance(expr, ColExpr) else expr + bound[name] = self._coerce_assign_operand(value) + + view = CTable._make_view(self, self._valid_rows) + view._computed_cols = dict(self._computed_cols) + view.col_names = list(self.col_names) + view._col_widths = dict(self._col_widths) + for name, value in bound.items(): + desc = view._normalize_transformer(value) + if desc["kind"] == "dsl": + kernel = desc["kernel"] + col_deps = desc["col_deps"] + view._computed_cols[name] = { + "kind": "dsl", + "dsl_source": kernel.dsl_source, + "kernel": kernel, + "col_deps": col_deps, + "dtype": view._dsl_result_dtype(kernel, col_deps, None), + "jit_backend": desc.get("jit_backend"), + } + else: + lazy = desc["lazy"] + view._computed_cols[name] = { + "kind": "expression", + "expression": lazy.expression, + "col_deps": desc["col_deps"], + "lazy": lazy, + "dtype": lazy.dtype, + } + view.col_names.append(name) + view._col_widths[name] = max(len(name), 15) + return view + # ------------------------------------------------------------------ # Column / row access # ------------------------------------------------------------------ @@ -10393,6 +10619,8 @@ def __getitem__(self, key): the base afterwards, and may still reference rows later deleted from the base). """ + if isinstance(key, ColExpr): + key = key._bind(self) if isinstance(key, str): physical = self._logical_to_physical_name(key) if physical in self._cols or physical in self._computed_cols: @@ -12207,7 +12435,7 @@ def dropna(self, subset: list[str] | None = None) -> CTable: def where( # noqa: C901 self, - expr_result: str | np.ndarray | blosc2.NDArray | blosc2.LazyExpr | blosc2.LazyUDF | Column, + expr_result: str | np.ndarray | blosc2.NDArray | blosc2.LazyExpr | blosc2.LazyUDF | Column | ColExpr, *, columns: list[str] | tuple[str, ...] | None = None, ) -> CTable: @@ -12219,10 +12447,10 @@ def where( # noqa: C901 The predicate can be supplied as a boolean :class:`blosc2.LazyExpr`, a boolean :class:`blosc2.NDArray`, a boolean NumPy array, a boolean - ``Column``, a :class:`blosc2.LazyUDF` (including those backed by a - :func:`blosc2.dsl_kernel`), or a string expression evaluated against - this table's columns. String expressions can reference stored and - computed columns directly by name. + ``Column``, an unbound :func:`col` expression, a :class:`blosc2.LazyUDF` + (including those backed by a :func:`blosc2.dsl_kernel`), or a string + expression evaluated against this table's columns. String expressions + can reference stored and computed columns directly by name. The returned object is a :class:`CTable` view sharing the original column data. The row-selection mask is evaluated immediately and @@ -12295,6 +12523,8 @@ def where( # noqa: C901 t.where((t.x > 0) and (t.y < 10)) t.where(not t.returned) """ + if isinstance(expr_result, ColExpr): + expr_result = expr_result._bind(self) if isinstance(expr_result, str): self._guard_varlen_scalar_expression(expr_result) operands = self._where_expression_operands(expr_result) diff --git a/tests/ctable/test_col_expr.py b/tests/ctable/test_col_expr.py new file mode 100644 index 000000000..04a35fd48 --- /dev/null +++ b/tests/ctable/test_col_expr.py @@ -0,0 +1,132 @@ +####################################################################### +# Copyright (c) 2019-present, Blosc Development Team +# All rights reserved. +# +# SPDX-License-Identifier: BSD-3-Clause +####################################################################### + +"""Tests for the unbound column expression (col()) and CTable.assign().""" + +from __future__ import annotations + +from dataclasses import dataclass + +import numpy as np +import pytest + +import blosc2 +from blosc2 import col +from blosc2.ctable import ColExpr + + +@dataclass +class Ledger: + revenue: float = blosc2.field(blosc2.float64()) + cost: float = blosc2.field(blosc2.float64()) + + +@dataclass +class Nullable: + x: float = blosc2.field(blosc2.float64(null_value=float("nan"))) + + +def _make_ledger(n: int = 5) -> blosc2.CTable: + data = [(float(10 * (i + 1)), float(i + 1)) for i in range(n)] + return blosc2.CTable(Ledger, data) + + +def test_assign_basic(): + t = _make_ledger() + t2 = t.assign(profit=col("revenue") - col("cost")) + assert t2.col_names == ["revenue", "cost", "profit"] + np.testing.assert_allclose(t2.profit[:], t.revenue[:] - t.cost[:]) + # original table is untouched + assert t.col_names == ["revenue", "cost"] + assert "profit" not in t._computed_cols + + +def test_assign_chain_end_to_end(): + t = _make_ledger() + result = ( + t.assign(profit=col("revenue") - col("cost"))[col("profit") > 0] + .sort_by("profit", ascending=False) + .head(10) + ) + pd = pytest.importorskip("pandas") + df = pd.DataFrame({"revenue": t.revenue[:], "cost": t.cost[:]}) + expected = ( + df.assign(profit=df.revenue - df.cost) + .query("profit > 0") + .sort_values("profit", ascending=False) + .head(10) + ) + np.testing.assert_allclose(result.profit[:], expected["profit"].to_numpy()) + + +def test_colexpr_filter_matches_bound_form(): + t = _make_ledger() + view_colexpr = t[col("revenue") > 25] + view_bound = t[t.revenue > 25] + np.testing.assert_array_equal(view_colexpr.revenue[:], view_bound.revenue[:]) + + +def test_assign_null_propagation(): + t = blosc2.CTable(Nullable, [(1.0,), (float("nan"),), (3.0,)]) + t2 = t.assign(y=col("x") + 1) + expected = (t.x + 1)[: t.nrows] # bound form, already tested elsewhere + np.testing.assert_array_equal(np.isnan(t2.y[:]), np.isnan(expected)) + np.testing.assert_allclose(t2.y[:][~np.isnan(t2.y[:])], expected[~np.isnan(expected)]) + + filtered_colexpr = t[col("x") < 0] + filtered_bound = t[t.x < 0] + assert filtered_colexpr.nrows == filtered_bound.nrows == 0 + + +def test_assign_reflected_scalar(): + t = _make_ledger() + t2 = t.assign(y=100 - col("revenue")) + np.testing.assert_allclose(t2.y[:], 100 - t.revenue[:]) + + +def test_col_unknown_name_fails_at_bind_time(): + expr = col("nope") # construction does not raise + assert isinstance(expr, ColExpr) + t = _make_ledger() + with pytest.raises(ValueError): + t.assign(z=expr) + + +def test_col_method_call_raises_clear_error(): + with pytest.raises(AttributeError, match="not supported on an unbound column expression"): + col("x").sum() + + +def test_colexpr_reused_across_tables(): + t1 = _make_ledger(3) + t2 = _make_ledger(4) + expr = col("revenue") + 1 + r1 = t1.assign(y=expr) + r2 = t2.assign(y=expr) + np.testing.assert_allclose(r1.y[:], t1.revenue[:] + 1) + np.testing.assert_allclose(r2.y[:], t2.revenue[:] + 1) + + +def test_assign_on_view(): + t = _make_ledger() + view = t[t.revenue > 20] + assigned = view.assign(z=col("revenue") * 2) + np.testing.assert_allclose(assigned.z[:], view.revenue[:] * 2) + + +def test_assign_result_is_read_only_view(): + t = _make_ledger() + t2 = t.assign(profit=col("revenue") - col("cost")) + assert t2._cols is t._cols + with pytest.raises(ValueError): + t2["revenue"][:] = np.zeros(t.nrows) + + +def test_assign_duplicate_name_raises(): + t = _make_ledger() + with pytest.raises(ValueError): + t.assign(revenue=col("cost")) From a41e21d8387bb3bf31fb692c61e2cd1a6997ffc5 Mon Sep 17 00:00:00 2001 From: Francesc Alted Date: Thu, 16 Jul 2026 12:27:27 +0200 Subject: [PATCH 3/8] P4: reject segmented dsl_kernel groupby UDF acceleration, keep bugfix Built and benchmarked the segmented reduceat-based UDF aggregation path specified for P4 (1e6 rows, 100k groups, a.max()-a.min()). Measured 1.2x against the plan's own >=5x gate (checked at other scales too: 1.1x-1.4x, never close), root-caused via cProfile to per-chunk-per-group Python bookkeeping in _merge_partials/_compute_partials that both the loop and segmented paths pay identically and that segmentation cannot skip. Per the plan's own benchmark-gate rule, the dispatch is not merged; the code was reverted rather than left in disabled. Kept the one fix worth keeping on its own merits: a @blosc2.dsl_kernel- decorated function passed as a groupby UDF aggregation crashed unconditionally (DSLKernel.__call__ uses a different calling convention than this call site expects) -- now unwrapped to the plain function. Co-Authored-By: Claude Sonnet 5 --- RELEASE_NOTES.md | 3 ++ plans/enhancing-ctable-phase2.md | 49 ++++++++++++++++++++++++++++++++ src/blosc2/groupby.py | 9 +++++- tests/ctable/test_groupby.py | 16 +++++++++++ 4 files changed, 76 insertions(+), 1 deletion(-) diff --git a/RELEASE_NOTES.md b/RELEASE_NOTES.md index bb7c95b5a..04541637d 100644 --- a/RELEASE_NOTES.md +++ b/RELEASE_NOTES.md @@ -15,6 +15,9 @@ XXX version-specific blurb XXX ### Bug fixes +- Fixed a `@blosc2.dsl_kernel`-decorated function crashing unconditionally + when passed as a groupby UDF aggregation (`g.agg(name=(col, + dsl_kernel_fn))`): it now runs like the equivalent undecorated callable. - Fixed `CTable.head()`/`tail()` silently discarding row order when called on a lazily-sorted view (e.g. `t.sort_by("col", ascending=False)` on a view, or any `.sort_by()` result chained off a prior filter): they diff --git a/plans/enhancing-ctable-phase2.md b/plans/enhancing-ctable-phase2.md index 7f883cdfa..48b968643 100644 --- a/plans/enhancing-ctable-phase2.md +++ b/plans/enhancing-ctable-phase2.md @@ -659,6 +659,55 @@ segmented numpy — no per-group Python. 3. High-cardinality smoke test: 10k groups, values equal between paths. 4. The benchmark script, recorded but not part of CI. +### Implementation notes (benchmark gate failed — NOT merged) + +Built the full segmented path as specified: `ast`-based recognition of the +whitelisted pattern (`_segmented_udf_plan`/`_SegmentedUDFTransformer` in +`groupby.py`, since deleted), replacing each whitelisted reduction call with +a placeholder bound to a `np.ufunc.reduceat` result and re-evaluating the +surrounding scalar arithmetic via a compiled expression: correct results, +verified against the loop path for every whitelisted reduction, the two +named composites, nulls, `chunk_size=2` chunk-straddling, and a 10k-group +smoke test. + +**Benchmark gate (plan's own spec: 1e6 rows, 100k groups, `a.max()-a.min()`) +measured 1.2x, not the required ≥5x** — checked at 5e6/500k and 1e6/500k +groups too (1.1x–1.4x, never close). Per the plan's own cross-cutting rule +4, **this means the segmented dispatch is not merged**; the code was +reverted rather than left in place disabled. + +Root cause (confirmed by `cProfile`, not guessed): the "1e5 Python calls to +a cheap UDF" cost the plan's estimate was built on is real but small next to +the *rest* of the group-by pipeline that both the loop path and the +segmented path pay identically and which this phase-1 architecture cannot +skip for exotic/multi-key group-by (`_factorize_keys`, `_compute_partials`, +and especially `_merge_partials`'s per-chunk-per-group Python-level +list-append bookkeeping that builds `_AggState.value`). Segmenting only +replaces the last step (call-the-UDF-per-group) with vectorized `reduceat`; +it does not and structurally cannot touch the bookkeeping steps before it, +which dominate wall time at every cardinality tested. A genuine 5x would +need bypassing that dict-of-Python-objects accumulator entirely (e.g. a +global vectorized sort-by-group-id across all chunks) — a materially larger +rewrite than "intercept in `_final_rows`," out of scope for this item as +specified. + +**What was kept**, because it is an independent, verified correctness fix +with no performance claim attached: a `@blosc2.dsl_kernel`-decorated +function passed as a groupby UDF aggregation (`g.agg(name=(col, +dsl_kernel_fn))`) previously crashed unconditionally — `DSLKernel.__call__` +uses the `(inputs_tuple, output, offset)` array-kernel convention, not the +"one array in, one scalar out" convention this call site uses, so +`spec.udf(group_values)` raised `TypeError: __call__() missing 1 required +positional argument: 'output'` wrapped in a `RuntimeError`, for *every* +group, regardless of what the kernel's body did. Fixed by calling +`spec.udf.func` (the plain wrapped function) when `spec.udf` is a +`DSLKernel` instance. Test: +`test_agg_udf_accepts_dsl_kernel_decorated_function` in +`tests/ctable/test_groupby.py`. This means a user who writes a groupby UDF +and later decorates it with `@blosc2.dsl_kernel` (e.g. to also use it +elsewhere as an elementwise kernel) no longer gets a crash — it just runs at +loop-path speed, same as an undecorated callable. + --- ## P5 — Mask-based nullable columns: PARKED (design recorded, do not build yet) diff --git a/src/blosc2/groupby.py b/src/blosc2/groupby.py index d760ec654..cce771e85 100644 --- a/src/blosc2/groupby.py +++ b/src/blosc2/groupby.py @@ -22,6 +22,7 @@ import numpy as np +from blosc2.dsl_kernel import DSLKernel from blosc2.schema import DictionarySpec, NDArraySpec, SchemaSpec, float64, int64 from blosc2.schema import bool as b2_bool from blosc2.schema import field as b2_field @@ -1804,8 +1805,14 @@ def _final_rows( # noqa: C901 row[spec.output_col] = _empty_udf_group else: group_values = np.concatenate(chunks) + # A @blosc2.dsl_kernel-decorated UDF is a DSLKernel + # instance whose __call__ expects the array-kernel + # calling convention (inputs_tuple, output, offset), + # not this "one array in, one scalar out" aggregation + # convention -- call the wrapped plain function instead. + udf_callable = spec.udf.func if isinstance(spec.udf, DSLKernel) else spec.udf try: - result = _python_scalar(spec.udf(group_values)) + result = _python_scalar(udf_callable(group_values)) except Exception as exc: raise RuntimeError( f"UDF aggregation {spec.output_col!r} raised for group " diff --git a/tests/ctable/test_groupby.py b/tests/ctable/test_groupby.py index cfd23ede5..3830d7efb 100644 --- a/tests/ctable/test_groupby.py +++ b/tests/ctable/test_groupby.py @@ -846,6 +846,22 @@ def test_agg_udf_matches_pandas_reference(): assert col(result, "city") == list(expected.index) +def test_agg_udf_accepts_dsl_kernel_decorated_function(): + """A @blosc2.dsl_kernel-decorated function is a DSLKernel instance whose + __call__ uses the array-kernel calling convention, not the "one array + in, one scalar out" convention this aggregation path calls with -- the + underlying plain function must be used instead.""" + + @blosc2.dsl_kernel + def udf_range(a): + return a.max() - a.min() + + t = CTable(SalesRow, new_data=DATA) + result = t.group_by("city", sort=True).agg(rng=("sales", udf_range)) + expected = t.group_by("city", sort=True).agg(rng=("sales", lambda a: a.max() - a.min())) + np.testing.assert_allclose(col(result, "rng"), col(expected, "rng")) + + def test_agg_udf_receives_only_live_nonnull_values(): seen = [] From cb166f89a204785a82535fae82d3a2484a85ea9f Mon Sep 17 00:00:00 2001 From: Francesc Alted Date: Thu, 16 Jul 2026 12:49:46 +0200 Subject: [PATCH 4/8] Add pandas engine example with measured timings Demonstrates engine=blosc2.jit for DataFrame.apply/Series.map, including the axis=0 speedup (measured, not just asserted) and the honest axis=1 result (no speedup, typically slower), plus the clear error for non-numeric columns. Co-Authored-By: Claude Sonnet 5 --- examples/pandas_engine.py | 135 ++++++++++++++++++++++++++++++++++++++ 1 file changed, 135 insertions(+) create mode 100644 examples/pandas_engine.py diff --git a/examples/pandas_engine.py b/examples/pandas_engine.py new file mode 100644 index 000000000..878bb6519 --- /dev/null +++ b/examples/pandas_engine.py @@ -0,0 +1,135 @@ +####################################################################### +# Copyright (c) 2019-present, Blosc Development Team +# All rights reserved. +# +# SPDX-License-Identifier: BSD-3-Clause +####################################################################### + +# Using blosc2.jit as a pandas execution engine. +# +# pandas' DataFrame.apply and Series.map accept an engine= argument: a +# callable exposing a __pandas_udf__ attribute that pandas dispatches to +# instead of running its own per-row/per-element Python loop. blosc2.jit is +# such an engine — but the function passed to it must be *vectorized*: it +# is called once with a full NumPy array (a column, or a row, depending on +# axis), not once per element. +# +# This example shows: +# 1. DataFrame.apply(f, engine=blosc2.jit) with axis=0 (the default) — the +# case that actually benefits from the engine, and why axis=1 does not. +# 2. Series.map(f, engine=blosc2.jit). +# 3. The clear error raised for non-numeric columns instead of a deep +# numexpr failure. +# 4. Timings for axis=0 (real speedup) vs axis=1 (no speedup) on a large +# DataFrame, so the difference is not just asserted but measured. + +try: + import pandas as pd +except ImportError: + raise SystemExit("This example requires pandas: pip install pandas") from None + +from time import perf_counter + +import numpy as np + +import blosc2 + +# --------------------------------------------------------------------------- +# 1. DataFrame.apply — axis=0 (column-wise) is where the engine wins +# --------------------------------------------------------------------------- + +df = pd.DataFrame( + { + "a": np.linspace(0, 10, 8), + "b": np.linspace(10, 20, 8), + } +) +print("Input DataFrame:") +print(df) + + +def transform(col): + return np.sin(col) * np.cos(col) + col**2 + + +expected = df.apply(transform) +result = df.apply(transform, engine=blosc2.jit) +print("\ndf.apply(transform, engine=blosc2.jit):") +print(result) +assert isinstance(result, pd.DataFrame) # not a raw ndarray +pd.testing.assert_frame_equal(result, expected) +print("matches plain df.apply(transform): True") + +# axis=1 (row-wise) still calls the function once per row, same as plain +# pandas — for a handful of columns the per-call wrapping overhead of the +# compute engine is *larger* than the win, so engine=blosc2.jit is typically +# slower than plain apply(axis=1) there. Prefer axis=0, as above. +result_axis1 = df.apply(transform, engine=blosc2.jit, axis=1) +pd.testing.assert_frame_equal(result_axis1, df.apply(transform, axis=1)) +print("\ndf.apply(transform, engine=blosc2.jit, axis=1) also matches (just not faster).") + +# --------------------------------------------------------------------------- +# 1b. Timings: axis=0 is a real win, axis=1 is not — measured, not asserted +# --------------------------------------------------------------------------- + +N_ROWS, N_COLS = 1_000_000, 8 +big_df = pd.DataFrame( + {f"c{i}": np.random.default_rng(i).random(N_ROWS) for i in range(N_COLS)}, +) + + +def heavier_transform(col): + return np.sin(col) * np.cos(col) + col**2 - np.sqrt(np.abs(col)) + np.exp(-col) + + +def timeit(fn, reps=3): + best = float("inf") + for _ in range(reps): + t0 = perf_counter() + fn() + best = min(best, perf_counter() - t0) + return best + + +print(f"\nTimings on a {N_ROWS:,}-row, {N_COLS}-column DataFrame (min of 3 runs):") + +t_plain0 = timeit(lambda: big_df.apply(heavier_transform)) +t_engine0 = timeit(lambda: big_df.apply(heavier_transform, engine=blosc2.jit)) +print(f" axis=0 plain df.apply(f): {t_plain0 * 1000:8.1f} ms") +print( + f" axis=0 df.apply(f, engine=blosc2.jit): {t_engine0 * 1000:7.1f} ms speedup: {t_plain0 / t_engine0:.1f}x" +) + +# axis=1 still calls the function once per *row* even with the engine, so +# it does not scale to 1e6 rows in a demo — a much smaller frame is enough +# to show the (lack of) speedup. +small_df = big_df.iloc[:20_000] +t_plain1 = timeit(lambda: small_df.apply(heavier_transform, axis=1)) +t_engine1 = timeit(lambda: small_df.apply(heavier_transform, engine=blosc2.jit, axis=1)) +print(f"\nSame comparison on axis=1, {len(small_df):,} rows (axis=1 doesn't scale to 1e6 in a demo):") +print(f" axis=1 plain df.apply(f, axis=1): {t_plain1 * 1000:8.1f} ms") +print( + f" axis=1 df.apply(f, engine=blosc2.jit, axis=1): {t_engine1 * 1000:7.1f} ms speedup: {t_plain1 / t_engine1:.2f}x" +) + +# --------------------------------------------------------------------------- +# 2. Series.map +# --------------------------------------------------------------------------- + +s = pd.Series(np.linspace(-5, 5, 6)) +mapped = s.map(lambda x: x**2 - 1, engine=blosc2.jit) +print("\nSeries.map(f, engine=blosc2.jit):") +print(mapped) +pd.testing.assert_series_equal(mapped, s.map(lambda x: x**2 - 1)) +print("matches plain Series.map: True") + +# --------------------------------------------------------------------------- +# 3. Non-numeric columns raise a clear error +# --------------------------------------------------------------------------- + +df_text = pd.DataFrame({"label": ["x", "y", "z"]}) +try: + df_text.apply(lambda x: x + 1, engine=blosc2.jit) +except ValueError as exc: + print("\nNon-numeric column raises ValueError:") + print(" ", exc) From c6e0658a113756d1b43f23a9743df3ed02e2f30d Mon Sep 17 00:00:00 2001 From: Francesc Alted Date: Thu, 16 Jul 2026 12:54:22 +0200 Subject: [PATCH 5/8] Drop axis=1 timing from pandas engine example, fix alignment axis=1 is slower with the engine (already noted in the code comment); showing timings for it added noise without a point. Also fixes column alignment in the axis=0 timing output. Co-Authored-By: Claude Sonnet 5 --- examples/pandas_engine.py | 21 ++++----------------- 1 file changed, 4 insertions(+), 17 deletions(-) diff --git a/examples/pandas_engine.py b/examples/pandas_engine.py index 878bb6519..a94917b43 100644 --- a/examples/pandas_engine.py +++ b/examples/pandas_engine.py @@ -20,8 +20,7 @@ # 2. Series.map(f, engine=blosc2.jit). # 3. The clear error raised for non-numeric columns instead of a deep # numexpr failure. -# 4. Timings for axis=0 (real speedup) vs axis=1 (no speedup) on a large -# DataFrame, so the difference is not just asserted but measured. +# 4. Measured timings for axis=0 on a large DataFrame — the actual win. try: import pandas as pd @@ -69,7 +68,7 @@ def transform(col): print("\ndf.apply(transform, engine=blosc2.jit, axis=1) also matches (just not faster).") # --------------------------------------------------------------------------- -# 1b. Timings: axis=0 is a real win, axis=1 is not — measured, not asserted +# 1b. Timings: axis=0 is a real, measured win # --------------------------------------------------------------------------- N_ROWS, N_COLS = 1_000_000, 8 @@ -95,21 +94,9 @@ def timeit(fn, reps=3): t_plain0 = timeit(lambda: big_df.apply(heavier_transform)) t_engine0 = timeit(lambda: big_df.apply(heavier_transform, engine=blosc2.jit)) -print(f" axis=0 plain df.apply(f): {t_plain0 * 1000:8.1f} ms") +print(f" {'plain df.apply(f):':<34s} {t_plain0 * 1000:7.1f} ms") print( - f" axis=0 df.apply(f, engine=blosc2.jit): {t_engine0 * 1000:7.1f} ms speedup: {t_plain0 / t_engine0:.1f}x" -) - -# axis=1 still calls the function once per *row* even with the engine, so -# it does not scale to 1e6 rows in a demo — a much smaller frame is enough -# to show the (lack of) speedup. -small_df = big_df.iloc[:20_000] -t_plain1 = timeit(lambda: small_df.apply(heavier_transform, axis=1)) -t_engine1 = timeit(lambda: small_df.apply(heavier_transform, engine=blosc2.jit, axis=1)) -print(f"\nSame comparison on axis=1, {len(small_df):,} rows (axis=1 doesn't scale to 1e6 in a demo):") -print(f" axis=1 plain df.apply(f, axis=1): {t_plain1 * 1000:8.1f} ms") -print( - f" axis=1 df.apply(f, engine=blosc2.jit, axis=1): {t_engine1 * 1000:7.1f} ms speedup: {t_plain1 / t_engine1:.2f}x" + f" {'df.apply(f, engine=blosc2.jit):':<34s} {t_engine0 * 1000:7.1f} ms speedup: {t_plain0 / t_engine0:.1f}x" ) # --------------------------------------------------------------------------- From e98c09f656670199d5ed57cb9122a041767deead Mon Sep 17 00:00:00 2001 From: Francesc Alted Date: Thu, 16 Jul 2026 12:55:59 +0200 Subject: [PATCH 6/8] Add CTable.assign()/blosc2.col() example Demonstrates the pandas-3 style chain, col() reuse across tables, null propagation, assign() on a view plus its read-only/no-copy guarantees, and the two documented error cases (unknown column at bind time, no method calls unbound). Co-Authored-By: Claude Sonnet 5 --- examples/ctable/assign_col.py | 149 ++++++++++++++++++++++++++++++++++ 1 file changed, 149 insertions(+) create mode 100644 examples/ctable/assign_col.py diff --git a/examples/ctable/assign_col.py b/examples/ctable/assign_col.py new file mode 100644 index 000000000..9d0b54e5e --- /dev/null +++ b/examples/ctable/assign_col.py @@ -0,0 +1,149 @@ +####################################################################### +# Copyright (c) 2019-present, Blosc Development Team +# All rights reserved. +# +# SPDX-License-Identifier: BSD-3-Clause +####################################################################### + +# CTable.assign() + blosc2.col(): pandas-3 style chaining. +# +# blosc2.col(name) builds an *unbound* column expression — a deferred +# recipe that only resolves against a table's columns once it is bound +# (passed to assign(), used to index/filter a table, or passed to +# where()). It reuses the exact same Column/NullableExpr operator +# machinery as the bound form (t.x + 1), so null propagation and +# comparison semantics are identical either way. +# +# CTable.assign(**named_exprs) returns a *view* with additional computed +# columns — it never mutates the table and never copies column data. +# +# This example shows: +# 1. The headline chain: assign -> filter -> sort -> head, in one line. +# 2. col() reused across different tables (it's just a recipe). +# 3. Null propagation rides along automatically. +# 4. assign() on a view, and the read-only guard on its result. +# 5. Errors: an unknown column name fails at bind time, not construction; +# method calls are not supported unbound. + +from dataclasses import dataclass + +import numpy as np + +import blosc2 +from blosc2 import col + +# --------------------------------------------------------------------------- +# Schema: a small ledger of revenue/cost per line item +# --------------------------------------------------------------------------- + + +@dataclass +class Ledger: + item: str = blosc2.field(blosc2.string(max_length=12)) + revenue: float = blosc2.field(blosc2.float64()) + cost: float = blosc2.field(blosc2.float64()) + + +LEDGER = [ + ("widgets", 100.0, 40.0), + ("gadgets", 50.0, 60.0), + ("gizmos", 30.0, 10.0), + ("doohickeys", 80.0, 90.0), + ("thingamajigs", 120.0, 45.0), +] + +t = blosc2.CTable(Ledger, new_data=LEDGER) +print("Ledger:") +print(t) + +# --------------------------------------------------------------------------- +# 1. The headline chain +# --------------------------------------------------------------------------- + +result = ( + t.assign(profit=col("revenue") - col("cost"))[col("profit") > 0] + .sort_by("profit", ascending=False) + .head(10) +) +print("\nProfitable items, most profitable first:") +print(result) + +# t.assign(...) never mutates t — it has no "profit" column +print(f"\nOriginal table columns (unchanged) : {t.col_names}") +print(f"Result table columns (with profit) : {result.col_names}") + +# --------------------------------------------------------------------------- +# 2. col() is a reusable recipe, not tied to any one table +# --------------------------------------------------------------------------- + +margin_expr = col("revenue") - col("cost") # built once + +t2 = blosc2.CTable(Ledger, new_data=[("extra_a", 200.0, 50.0), ("extra_b", 20.0, 25.0)]) + +print("\nThe same expression applied to two different tables:") +print(" ", t.assign(margin=margin_expr).margin[:]) +print(" ", t2.assign(margin=margin_expr).margin[:]) + +# Reflected operands and constants work as expected +scaled = t.assign(revenue_pct=100 * col("revenue") / (col("revenue") + col("cost"))) +print("\nRevenue share of (revenue + cost), reflected/scalar operands:") +print(" ", np.round(scaled.revenue_pct[:], 1)) + +# --------------------------------------------------------------------------- +# 3. Null propagation rides along automatically +# --------------------------------------------------------------------------- + + +@dataclass +class LedgerNullable: + item: str = blosc2.field(blosc2.string(max_length=12)) + revenue: float = blosc2.field(blosc2.float64(null_value=float("nan"))) + cost: float = blosc2.field(blosc2.float64()) + + +nullable_data = [ + ("a", 100.0, 40.0), + ("b", float("nan"), 60.0), # missing revenue + ("c", 30.0, 10.0), +] +tn = blosc2.CTable(LedgerNullable, new_data=nullable_data) +tn_assigned = tn.assign(profit=col("revenue") - col("cost")) +print("\nNull propagation — item 'b' has no revenue, profit is null (NaN):") +print(tn_assigned.select(["item", "profit"])) + +# --------------------------------------------------------------------------- +# 4. assign() on a view, and the read-only guard on the result +# --------------------------------------------------------------------------- + +view = t[t.revenue > 60] +view_assigned = view.assign(margin=col("revenue") - col("cost")) +print("\nassign() on a filtered view (revenue > 60):") +print(view_assigned) + +try: + view_assigned["revenue"][:] = np.zeros(len(view_assigned)) +except ValueError as exc: + print(f"\nWrite guard on assign() result: {exc}") + +# assign() shares column storage — no copy of the underlying data +same_storage = view_assigned._cols is t._cols +print(f"\nassign() shares column storage with the base table: {same_storage}") + +# --------------------------------------------------------------------------- +# 5. Errors: unknown columns fail at bind time; no method calls unbound +# --------------------------------------------------------------------------- + +expr = col("nope") # constructing col() never raises +print("\ncol('nope') constructed fine; error only appears once bound:") +try: + t.assign(x=expr) +except ValueError as exc: + print(" ", exc) + +try: + col("revenue").sum() +except AttributeError as exc: + print("\nMethod calls are not supported on an unbound column expression:") + print(" ", exc) + +print("\nDone.") From 29ffd46295151fa6cc2a555e014e6ae08e7a74d5 Mon Sep 17 00:00:00 2001 From: Francesc Alted Date: Thu, 16 Jul 2026 13:07:16 +0200 Subject: [PATCH 7/8] Added an example that bundles ctables and ndarrays in the same TreeStore --- examples/ctable/treestore_bundle.py | 141 ++++++++++++++++++++++++++++ 1 file changed, 141 insertions(+) create mode 100644 examples/ctable/treestore_bundle.py diff --git a/examples/ctable/treestore_bundle.py b/examples/ctable/treestore_bundle.py new file mode 100644 index 000000000..7c8a279f5 --- /dev/null +++ b/examples/ctable/treestore_bundle.py @@ -0,0 +1,141 @@ +####################################################################### +# Copyright (c) 2019-present, Blosc Development Team +# All rights reserved. +# +# SPDX-License-Identifier: BSD-3-Clause +####################################################################### + +# Bundling CTables and NDArrays together in a single TreeStore file. +# +# A TreeStore can hold both plain NDArrays and CTable objects in the same +# .b2d directory or .b2z zip archive. Each CTable is stored inline as a +# named subtree — its columns, metadata, and index sidecars all live as +# normal Blosc2 leaves inside the bundle. From the outside the CTable +# appears as a single key, just like any other leaf. + +import shutil +import tempfile +from dataclasses import dataclass + +import numpy as np + +import blosc2 + + +@dataclass +class Sensor: + sensor_id: int = blosc2.field(blosc2.int32(ge=0)) + temperature: float = blosc2.field(blosc2.float64(), default=0.0) + day: int = blosc2.field(blosc2.int16(ge=1, le=365), default=1) + + +@dataclass +class Event: + event_id: int = 0 + label: str = "" + + +tmpdir = tempfile.mkdtemp(prefix="blosc2_bundle_") +b2d_path = f"{tmpdir}/dataset.b2d" +b2z_path = f"{tmpdir}/dataset.b2z" + +try: + # ------------------------------------------------------------------ + # 1. Build some in-memory data + # ------------------------------------------------------------------ + rng = np.random.default_rng(42) + N = 500 + + sensors = blosc2.CTable(Sensor) + for _ in range(N): + sensors.append( + Sensor( + sensor_id=int(rng.integers(0, 5)), + temperature=float(rng.normal(20.0, 5.0)), + day=int(rng.integers(1, 366)), + ) + ) + + events = blosc2.CTable(Event) + for i in range(20): + events.append(Event(event_id=i, label=f"evt-{i:03d}")) + + calibration = blosc2.linspace(0.0, 1.0, 10) + timestamps = np.arange(N, dtype=np.int64) + + # ------------------------------------------------------------------ + # 2. Write a .b2d bundle mixing CTables and NDArrays + # ------------------------------------------------------------------ + with blosc2.TreeStore(b2d_path, mode="w") as ts: + # NDArrays sit alongside CTables — same assignment syntax + ts["/raw/calibration"] = calibration + ts["/raw/timestamps"] = timestamps + + # CTables are stored inline; internals are hidden from normal traversal + ts["/tables/sensors"] = sensors + ts["/tables/events"] = events + + print("Keys written to .b2d bundle:") + for k in sorted(ts.keys()): + print(f" {k}") + + # ------------------------------------------------------------------ + # 3. Read the .b2d bundle back + # ------------------------------------------------------------------ + print("\nReading .b2d bundle:") + with blosc2.open(b2d_path, mode="r") as ts: + # CTable is returned transparently + s = ts["/tables/sensors"] + print(f" /tables/sensors : {type(s).__name__}, {len(s):,} rows") + print(f" mean temp : {s['temperature'].mean():.2f}") + + e = ts["/tables/events"] + print(f" /tables/events : {type(e).__name__}, {len(e)} rows") + + cal = ts["/raw/calibration"] + print(f" /raw/calibration : {type(cal).__name__}, shape={cal.shape}") + + # Object internals are not exposed in keys() + assert "/tables/sensors/_meta" not in ts + assert "/tables/sensors/_cols" not in ts + + # ------------------------------------------------------------------ + # 4. Pack the .b2d bundle into a single .b2z archive + # ------------------------------------------------------------------ + with blosc2.TreeStore(b2d_path, mode="r") as ts: + ts.to_b2z(filename=b2z_path) + + print(f"\nPacked to .b2z: {b2z_path}") + + # ------------------------------------------------------------------ + # 5. Read directly from the .b2z archive (offset-based, no extraction) + # ------------------------------------------------------------------ + print("\nReading .b2z archive:") + with blosc2.open(b2z_path, mode="r") as ts: + s2 = ts["/tables/sensors"] + print(f" /tables/sensors : {type(s2).__name__}, {len(s2):,} rows") + + e2 = ts["/tables/events"] + print(f" /tables/events : {type(e2).__name__}, {len(e2)} rows") + print(f" first event : event_id={e2[0]['event_id']}, label='{e2[0]['label']}'") + + # ------------------------------------------------------------------ + # 6. Append a new row to a CTable inside the bundle + # ------------------------------------------------------------------ + with blosc2.TreeStore(b2d_path, mode="a") as ts: + s3 = ts["/tables/sensors"] + s3.append(Sensor(sensor_id=99, temperature=-10.0, day=1)) + print(f"\nAfter append: /tables/sensors has {len(s3):,} rows") + # Closing the table explicitly is optional; TreeStore.__exit__ handles it + s3.close() + + # ------------------------------------------------------------------ + # 7. Delete a CTable from the bundle + # ------------------------------------------------------------------ + with blosc2.TreeStore(b2d_path, mode="a") as ts: + del ts["/tables/events"] + print(f"\nAfter deleting /tables/events, keys: {sorted(ts.keys())}") + +finally: + shutil.rmtree(tmpdir) + print("\nTemporary files removed.") From b8c56f6d2b89a3a626a3134ad0a8eb4f9eb4cab7 Mon Sep 17 00:00:00 2001 From: Francesc Alted Date: Thu, 16 Jul 2026 13:32:06 +0200 Subject: [PATCH 8/8] Review pass: docs for assign()/col(), robust pandas version check, phase-3 plan - Add CTable.assign and blosc2.col to the API reference, with a short chained-pipelines section. - Fix the assign() docstring: values are bound against self, not the view. - Make the pandas major-version test guard numeric instead of a string comparison. - Close the phase-2 plan and carry deferred P3/P5 forward into plans/enhancing-ctable-phase3.md. Co-Authored-By: Claude Fable 5 --- doc/reference/ctable.rst | 21 +++ plans/enhancing-ctable-phase2.md | 6 +- plans/enhancing-ctable-phase3.md | 276 +++++++++++++++++++++++++++++++ src/blosc2/ctable.py | 9 +- tests/test_pandas_udf_engine.py | 2 +- 5 files changed, 307 insertions(+), 7 deletions(-) create mode 100644 plans/enhancing-ctable-phase3.md diff --git a/doc/reference/ctable.rst b/doc/reference/ctable.rst index 67727e3bb..c875f3ec3 100644 --- a/doc/reference/ctable.rst +++ b/doc/reference/ctable.rst @@ -309,6 +309,23 @@ When a NumPy structured array is needed, materialize explicitly:: np.asarray(t[:10]) +Chained pipelines +~~~~~~~~~~~~~~~~~ + +:meth:`CTable.assign` returns a view with additional computed columns — +never mutating the table, never copying column data — and :func:`blosc2.col` +builds an unbound column expression that resolves against a table only when +bound (in ``assign()``, ``t[...]``, or :meth:`CTable.where`). Together they +enable pandas-3 style method chains:: + + from blosc2 import col + + result = ( + t.assign(profit=col("revenue") - col("cost"))[col("profit") > 0] + .sort_by("profit", ascending=False) + .head(10) + ) + .. autosummary:: CTable.where @@ -316,24 +333,28 @@ When a NumPy structured array is needed, materialize explicitly:: CTable.view CTable.take CTable.select + CTable.assign CTable.head CTable.tail CTable.sample CTable.sort_by CTable.iter_sorted CTable.group_by + col .. automethod:: CTable.where .. automethod:: CTable.dropna .. automethod:: CTable.view .. automethod:: CTable.take .. automethod:: CTable.select +.. automethod:: CTable.assign .. automethod:: CTable.head .. automethod:: CTable.tail .. automethod:: CTable.sample .. automethod:: CTable.sort_by .. automethod:: CTable.iter_sorted .. automethod:: CTable.group_by +.. autofunction:: col Group-by reductions diff --git a/plans/enhancing-ctable-phase2.md b/plans/enhancing-ctable-phase2.md index 48b968643..fe7f98701 100644 --- a/plans/enhancing-ctable-phase2.md +++ b/plans/enhancing-ctable-phase2.md @@ -1,6 +1,10 @@ # Enhancing CTable, phase 2: finishing the pandas-3 story -**Status:** NOT STARTED (plan written 2026-07-16). Successor to +**Status:** CLOSED (2026-07-16, branch `enhancing-ctable2`). P1 and P2 +landed; P4 was built, failed its own benchmark gate (1.2x vs required ≥5x) +and was deliberately not merged (its groupby-UDF crash fix was kept) — see +each item's "Implementation notes". P3 and P5 were deferred, carried +forward to `plans/enhancing-ctable-phase3.md`. Successor to `plans/enhancing-ctable.md` (phase 1, fully landed on branch `enhancing-ctable`: Arrow PyCapsule interchange, read-only views, the sentinel-null story including null propagation and `NullableExpr` diff --git a/plans/enhancing-ctable-phase3.md b/plans/enhancing-ctable-phase3.md new file mode 100644 index 000000000..9551b15e6 --- /dev/null +++ b/plans/enhancing-ctable-phase3.md @@ -0,0 +1,276 @@ +# Enhancing CTable, phase 3: variable-length strings and the null-mask question + +**Status:** NOT STARTED (plan written 2026-07-16). Successor to +`plans/enhancing-ctable-phase2.md` (phase 2, landed on branch +`enhancing-ctable2`: the pandas `engine=blosc2.jit` fix + `Series.map`, +`CTable.assign()` + unbound `blosc2.col()`, the lazily-sorted-view +`head()`/`tail()` ordering fix, and the DSLKernel groupby-UDF crash fix; +phase 2's P4 segmented-UDF fast path was built, failed its own benchmark +gate at 1.2x vs the required ≥5x, and was deliberately NOT merged — see +that plan's P4 implementation notes before ever reattempting it). + +**Audience:** an implementing model/developer who has NOT read the +discussions that produced this plan and has NOT read phases 1–2. Everything +needed is in this file. When in doubt, prefer the laziest change that +satisfies the acceptance criteria — do not add abstractions, protocols, or +options this plan does not ask for. + +**Important practical notes (carried over, still true):** + +- Line numbers below were verified on 2026-07-16 against branch + `enhancing-ctable2`. Lines WILL drift; always locate code by the symbol + names given, using grep, and treat line numbers only as hints. +- Run Python/pytest through the `blosc2` conda env: + `conda run -n blosc2 python -m pytest ...`. Never use the repo `.venv` + (stale). +- Editing `.pyx` files does NOT trigger a rebuild in an editable install; + prefer pure-Python implementations. +- **Docstrings and code comments must be self-contained.** Never reference + this plan, earlier phases, or item labels ("P3", "P3.d") from source, + tests, or bench scripts — state the semantics directly. This is an + explicit maintainer rule. +- CTable tests live under `tests/ctable/`; match the style of neighboring + tests. The dev env has pandas 3.0.3, numpy 2.4.6, pyarrow, duckdb, and + polars installed, so `importorskip` tests run for real there. +- The dev machine is an Apple-silicon Mac; benchmark targets below assume + it. + +--- + +## Scope decision (made 2026-07-16, do not re-litigate the ordering) + +Phase 2 deliberately shipped without its P3 (variable-length strings) and +P5 (mask-based nullable columns). The ordering question — "should the +mask-based null design be built first, since strings are the dtype where a +reserved sentinel is most awkward?" — was considered and answered: + +- **P3.a (storage + roundtrip) goes first.** It is null-representation- + agnostic (companion offsets+bytes arrays), it is the largest de-risking + step, and nothing in it forecloses either null design. +- **The sentinel-vs-mask call for utf8 is made at P3.b/P3.c time**, with + real data in hand. The default remains sentinel (consistent with every + other dtype). But free-text is the one dtype where *any* value is legal, + so a sentinel can collide with real data; if that proves lossy in + practice during Arrow interop (P3.b), that is a legitimate trigger for + P5's third unpark criterion — and the response is to build the mask + machinery *scoped to what utf8 needs*, not the full every-dtype P5. +- **P5 is otherwise still parked.** None of its unpark criteria have been + hit as of 2026-07-16. + +--- + +## P3 — First-class variable-length string columns (the real project) + +### Goal + +pandas 3's headline dtype is an efficient variable-length string with NA +semantics as the default string story. CTable's equivalent must make this +work, at full speed, with bounded memory: + +```python +@dataclass +class Row: + name: str = blosc2.field(blosc2.utf8()) # new: varlen, first-class + ... + + +t[t.name == "Paris"] # vectorized comparison +t.group_by("name").sum("x") # groupby key +t.sort_by("name") # ordering +``` + +### Why the existing pieces don't cover it (verified 2026-07-16) + +- `blosc2.string(max_length=n)` is fixed-width `U` (UTF-32): 4 bytes per + character per row before compression, and 32-byte comparisons that made + string groupby 8x slower than pandas until phase 1's hash fix. Wrong + answer for long/variable text. +- `blosc2.vlstring()` (`schema.py:762`) stores msgpack-serialized cells — + per-cell decode, no vectorized anything: `_ensure_queryable` + (`ctable.py:1715`) rejects comparisons, `groupby.py:139` rejects keys, + `ctable.py:10611` rejects sort (grep for these guards by message text, + the line numbers have drifted). Retrofitting msgpack cells to be fast is + a dead end; do not try. +- `blosc2.dictionary()` is the right tool for LOW-cardinality strings and + is already first-class. P3 is for high-cardinality/free-text columns. + +### Decisions already made + +1. **Storage layout: Arrow-style offsets + bytes.** Two companion NDArrays + per column in the store: `int64` offsets (length `n+1`) and a `uint8` + UTF-8 byte blob. Precedent for companion arrays already exists + (`valid_rows` in `ctable_storage.py`; the dictionary store; the phase-1 + decision that companion arrays need no format bump). Chunk-aligned + access: reading rows `[a, b)` needs `offsets[a : b+1]` plus + `bytes[offsets[a] : offsets[b]]` — both are plain NDArray slice reads. +2. **In-memory representation: `numpy.dtypes.StringDType`** (numpy ≥ 2.0; + verify the package's minimum numpy before relying on it; if the floor + is < 2.0, gate the feature on the installed numpy and raise a clear + error). StringDType arrays support `==`, ordering, and the vectorized + `np.strings` functions. +3. **Expression routing: chunked numpy, NOT numexpr.** numexpr/miniexpr + cannot evaluate StringDType. String comparisons must evaluate chunk by + chunk in numpy and produce the same physical-length boolean masks the + existing predicates produce. Look at how dictionary columns solved the + identical problem (`Column._dictionary_eq`, and grep how its + physical-slot predicates flow into `where()`) — mirror that pattern, do + not invent a new one. +4. **Nulls: sentinel-based by default**, consistent with everything else + (the existing machinery already supports string sentinels; grep + `test_null_value_string`) — **but this is the one decision with a + planned checkpoint**: see "Scope decision" above. If during P3.b/P3.c + the sentinel choice for free-text proves lossy against real Arrow data, + switch utf8 (and only utf8) to a companion validity-mask array, using + the P5 design below. Record whichever way it goes in the implementation + notes. +5. **Spec name `blosc2.utf8(nullable=..., null_value=...)`.** Keep + `string()` (fixed) and `vlstring()` (msgpack) untouched for backward + compatibility. Making `utf8` the default for `str` fields is explicitly + NOT part of this item — propose it separately once P3 has soaked. + +### Phasing (each lands separately, in order) + +**P3.a Storage + roundtrip.** Schema spec `utf8` in `schema.py` (mirror how +`vlstring` declares itself, but with `kind: "utf8"`); read/write paths in +`ctable_storage.py` for the two companion arrays; `append`/`extend`/ +`__getitem__` on the column returning StringDType arrays; persistence +roundtrip (`.b2z` save/open). In-place cell UPDATE of a varlen value +changes the byte length — decision: rewriting a cell rewrites the column's +tail (offsets shift). That is O(n) and fine for v1; `Column.__setitem__` on +a utf8 column should work but the docstring must state the cost. Tests: +roundtrip (ASCII, non-ASCII, empty strings, 1-char, multi-KB values), +append/extend, setitem, persistence, repr. + +**P3.b Arrow interop.** `iter_arrow_batches` exports utf8 columns as +`pa.large_string()` (always large — no int32-offset special-casing); +`from_arrow` maps incoming `string`/`large_string` columns to `utf8` specs +(today they land as fixed-width or dictionary — check `_auto_null_sentinel` +/ `from_arrow`'s type dispatch before changing it; keep the old mapping +available via the existing import options if one exists). Null cells map +sentinel↔validity like every other dtype — **this is the sentinel-vs-mask +checkpoint from decision 4**. Tests: `pa.table(ct)` roundtrip, duckdb query +on a utf8 column, `from_arrow(pa_table)` ingest. + +**P3.c Filters and expressions.** Carve utf8 out of the +`_ensure_queryable` rejection for comparisons only; implement `==`, `!=`, +`<`, `<=`, `>`, `>=` against scalars and other utf8 columns via the +chunked-numpy predicate path from decision 3, with the phase-1 SQL null +rule (a null never satisfies any comparison; grep `_null_aware_compare` +for the semantics to match). `blosc2.startswith`-style function ops already +work — add tests pinning them. Tests: filter correctness incl. null rows, +`t[t.name == "x"]` on views, comparison with non-string scalar raises +clearly. + +**P3.d Groupby keys.** Replace the groupby rejection for utf8 keys +(grep `Cannot group by variable-length` in `groupby.py`). Factorization: +read the chunk as offsets+bytes and hash rows vectorized — same trick as +phase 1's `_factorize_fixed_width_str` (`groupby.py`, study it first) but +over variable-length bytes: a vectorized loop over the (few) distinct +byte-lengths, or `np.frombuffer`-based chunked mixing; verify + +`np.unique`-on-StringDType fallback for collisions, identical output +contract. Benchmark against dictionary-key groupby on the same data +(`bench/ctable/bench_groupby_keys.py` — extend it); target: within 3x of +the dictionary path for 1e7 rows/low cardinality. If the vectorized hash +proves hard, the honest fallback is `np.unique` on the StringDType chunk +(correct, slower) — land correctness first, speed second. + +**P3.e Sort.** Lift the sort rejection (grep the `sort_by` varlen guard in +`ctable.py`): `sort_by` on utf8 uses `np.argsort` on the StringDType array +(chunked merge if the existing sort machinery is chunked — study +`sort_by`'s path first). Null ordering must match the existing convention +(nulls last; grep `test_sort_nulls_last`). + +**Out of scope for P3:** making `utf8` the `str`-field default; `.str` +accessor namespaces; regex operations beyond what `np.strings` gives; +string interning/dedup (that's `dictionary()`). + +Acceptance for the item overall: the three Goal-section lines work, the +groupby bench number is recorded, and `vlstring`/`string` behavior is +byte-for-byte unchanged. + +### Benchmark gate reminder + +Phase 1 rejected its own JIT groupby engine and phase 2 rejected its own +segmented-UDF path on measured numbers. P3.d has the same rule: do not +merge a "fast" factorization the recorded benchmark shows is not fast — +land the correct `np.unique` fallback instead and record the gap. + +--- + +## P5 — Mask-based nullable columns: STILL PARKED (design recorded, do not build yet) + +Sentinels cannot represent: nullable bool with all 256 byte values in use +(current nullable bool reserves 255), full-range `uint8`/`int8`, and any +dtype where reserving a value is unacceptable. The agreed design, recorded +in phase 1 and restated here so it survives: + +- A hidden companion boolean validity array per column — just another key + in the `.b2z` store, exactly like `valid_rows` (`ctable_storage.py`, + grep `valid_rows`) and P3's offsets array. **No .b2z or C-Blosc2 format + change.** +- A per-column schema marker (e.g. `null_mask: true` in the column's + metadata dict) — NOT a global `/_meta` `version` bump, so only tables + actually using masks are unreadable by old readers, failing cleanly at + schema load. +- Integration points when built: `Column.is_null()` reads the mask; + write paths (`append`/`extend`/`__setitem__`/`Column.assign`) maintain + it; `_nonnull_chunks` and the lazy reduction masks consult it; groupby + null handling; Arrow import/export maps mask↔validity directly (cheaper + than sentinel conversion); `fillna`/`dropna`. + +**Criteria to unpark** (any one): + +1. A user asks for nullable bool without the 255 reservation. +2. A user needs full-range small ints with nulls. +3. Arrow ingest of a type whose sentinel choice is provably lossy — **P3.b + is the most likely place this fires** (free-text utf8 is the dtype + where any sentinel value can collide with real data). If it fires + there, build the mask machinery scoped to utf8 only, behind the + per-column `null_mask` marker, and leave every other dtype on + sentinels. + +Until one fires, build nothing — every phase-1 and phase-2 feature works +on sentinels. + +--- + +## Small known gap (candidate side-item, not scheduled) + +Computed columns carry no null metadata: `t.add_computed_column("y", +"x + 1")` on a nullable `x` produces correct NaN propagation in the +*values*, but `t.y.is_null()` returns all-False (verified 2026-07-16; the +same holds for `CTable.assign`, which shares the machinery — this +predates `assign()`). The values are right; only the null *introspection* +on the derived column is blind. If P3 or a user bumps into this, the fix +belongs in the computed-column metadata (`_computed_cols` entries record a +dtype but no null sentinel); derive the sentinel from the expression's +NullableExpr provenance when available. Cheap to do alongside P3.c's null +comparison work; do not start it standalone without a use case. + +--- + +## Cross-cutting rules for the implementer (unchanged from phase 2) + +1. **Verify before building.** This codebase has been ahead of every + analysis so far. Before implementing any sub-item, grep for it; if it + exists, write the test that proves it and move on. +2. **One PR per sub-item** (P3.a … P3.e each land separately). P5: no PR + unless an unpark criterion fires. +3. **No new dependencies.** pandas/pyarrow/duckdb/polars appear only in + tests via `importorskip`; numpy StringDType is stdlib-of-the-project. +4. **Benchmark gates are real.** P3.d must not merge a "fast path" that + the recorded benchmark shows is not fast. Phases 1 and 2 each rejected + their own fast path on exactly this rule and were right to. +5. **Error messages name the escape hatch** (the view-write error pointing + at `take()`/`copy()` is the house style). +6. **Docstrings are self-contained** — no references to this plan or its + item labels anywhere in `src/`, `tests/`, or `bench/`. +7. **Docstrings in the existing style** (NumPy-doc with Examples sections, + as throughout `ctable.py`). +8. Run the CTable subset with + `conda run -n blosc2 python -m pytest tests/ctable -q`; run + `tests/ndarray` too when touching anything under `src/blosc2/` outside + `ctable.py`/`groupby.py`. +9. After landing a sub-item, append an "Implementation notes" subsection + to its section in THIS file: what landed, what deviated and why, + measured numbers, and commit hashes. diff --git a/src/blosc2/ctable.py b/src/blosc2/ctable.py index a3da2b840..57b8ecf1a 100644 --- a/src/blosc2/ctable.py +++ b/src/blosc2/ctable.py @@ -10368,11 +10368,10 @@ def assign(self, **named_exprs) -> CTable: .head(10) ) - Values are bound against the *new view* being built, one keyword at a - time, but against a snapshot taken before any of this call's new - columns were added — so a later keyword cannot reference an earlier - one from the same ``assign()`` call (that raises the usual unknown- - column error). Chain two ``assign()`` calls for that:: + Values are bound against ``self``, before any of this call's new + columns exist — so a later keyword cannot reference an earlier one + from the same ``assign()`` call (that raises the usual unknown-column + error). Chain two ``assign()`` calls for that:: t2 = t.assign(profit=col("revenue") - col("cost")) t3 = t2.assign(margin=col("profit") / col("revenue")) diff --git a/tests/test_pandas_udf_engine.py b/tests/test_pandas_udf_engine.py index cc576c996..7441345f9 100644 --- a/tests/test_pandas_udf_engine.py +++ b/tests/test_pandas_udf_engine.py @@ -138,7 +138,7 @@ def add_one(x): try: import pandas as pd - _pandas_too_old = pd.__version__ < "3" + _pandas_too_old = int(pd.__version__.split(".")[0]) < 3 except ImportError: pd = None _pandas_too_old = False