Skip to content
Merged
29 changes: 29 additions & 0 deletions RELEASE_NOTES.md
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,35 @@

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 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
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
`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
Expand Down
73 changes: 73 additions & 0 deletions bench/bench_pandas_engine.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,73 @@
#######################################################################
# Copyright (c) 2019-present, Blosc Development Team <blosc@blosc.org>
# 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()
1 change: 1 addition & 0 deletions doc/guides/index.rst
Original file line number Diff line number Diff line change
Expand Up @@ -13,6 +13,7 @@ Topics

optimization_tips
sharing_across_processes
pandas_engine

Command line tools
------------------
Expand Down
83 changes: 83 additions & 0 deletions doc/guides/pandas_engine.md
Original file line number Diff line number Diff line change
@@ -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
```
21 changes: 21 additions & 0 deletions doc/reference/ctable.rst
Original file line number Diff line number Diff line change
Expand Up @@ -309,31 +309,52 @@ 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
CTable.dropna
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
Expand Down
Loading
Loading