Skip to content

Commit 0af2a50

Browse files
timsaucerclaude
andcommitted
feat: let extension bundles declare scalar, aggregate, and window functions
`SessionExtensionComponents` gains `udfs`, `udafs`, and `udwfs`, so a library shipping functions can be installed with one `with_extensions` call instead of documenting a per-function `register_*` recipe. Either the Python wrapper or a raw capsule exportable is accepted; the registered name comes off the function. Installation now splits into a fallible part and an infallible one. Collecting hooks, building the codec chains, resolving the declared functions, and running the planner hooks all write nothing; only the final step binds the planner and registers. That keeps "nothing is written until every hook has returned" true now that components reach the shared `SessionState`, where there is nothing to roll back to. A new comment states the rule for whoever adds the next field. Two extensions declaring one name in a single call is a `ValueError` naming both, since a function registry has no fall-through the way a codec chain does. Shadowing a name the session already has stays legal, which `enable_spark_functions` relies on. `__post_init__` now normalizes fields by metadata rather than by the `_codecs` name suffix, so the new fields are covered and later ones will be too. `MyFunctionExtension` in `datafusion-ffi-example` declares this crate's three functions across a real FFI boundary. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
1 parent 1d629a9 commit 0af2a50

14 files changed

Lines changed: 758 additions & 85 deletions

File tree

‎docs/source/extension-guide/bundles.md‎

Lines changed: 55 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -21,7 +21,7 @@
2121

2222
# Extension bundles
2323

24-
If your library ships codecs, or a query planner, or both, expose a **bundle**
24+
If your library ships codecs, functions, or a query planner, expose a **bundle**
2525
and let callers install it with
2626
{py:meth}`~datafusion.SessionContext.with_extensions`. This is the recommended
2727
way to package an extension, and the rest of this page explains what the
@@ -45,6 +45,7 @@ class MyEngineExtension:
4545
return SessionExtensionComponents(
4646
logical_extension_codecs=(self._make_logical_codec(ctx),),
4747
physical_extension_codecs=(self._make_physical_codec(ctx),),
48+
udfs=(MyScalarUDF(),),
4849
)
4950

5051
def __datafusion_session_planner__(self, ctx: SessionContext, fallback):
@@ -55,15 +56,22 @@ class MyEngineExtension:
5556
```
5657

5758
Implement whichever apply: a codec-only library defines the first, a library
58-
that ships only an optimizing planner defines the second. The caller then
59-
writes:
59+
that ships only an optimizing planner defines the second, and a library that
60+
ships only functions defines the first and leaves the codec fields empty. The
61+
caller then writes:
6062

6163
```python
6264
ctx = SessionContext(config).with_extensions(lib_a.Extension(), lib_b.Extension())
6365
ctx.register_table("t", lib_a.TableProvider())
64-
ctx.register_udf(udf(lib_b.SomeUDF()))
6566
```
6667

68+
Declare functions rather than registering them yourself inside the hook.
69+
Declared components are resolved before anything is written, and they are
70+
registered after every codec is installed; a registration you make during the
71+
hook happens too early to see the other bundles' codecs and is not undone if a
72+
later extension fails. Table providers are still registered by the caller, on
73+
the returned handle — see {ref}`extension_bundles_transaction`.
74+
6775
`MyPlannerExtension` in [`datafusion-ffi-query-planner-example`] is a complete
6876
Rust implementation of the protocol, including taking the task-context provider
6977
off the supplied context, wrapping its codecs in `BundledLogicalCodec` /
@@ -281,14 +289,52 @@ for direct ones. The wrapper travels with the codec; the bundle does not.
281289
The query planner is exempt — it carries no wire id, so it may be an object or
282290
a capsule.
283291

292+
(extension_bundles_transaction)=
293+
284294
## Failure and rollback
285295

286296
Nothing is written to the session until every factory has returned and every
287-
capsule has been validated, so a factory that raises leaves the session exactly
288-
as it was. A factory that mutates the context it is handed — registering a
289-
table, say — is **not** rolled back, which is why bundle objects must be
290-
configuration-only: create fresh components on each call, never cache bound
291-
components, and do not retain the context passed in.
297+
component has been validated, so a factory that raises leaves the session
298+
exactly as it was. A factory that mutates the context it is handed —
299+
registering a table, say — is **not** rolled back, which is why bundle objects
300+
must be configuration-only: create fresh components on each call, never cache
301+
bound components, and do not retain the context passed in.
302+
303+
That guarantee is why the installation runs in the order it does. A call splits
304+
into a part that may fail and a part that may not:
305+
306+
1. **Collect.** Every `__datafusion_session_components__` runs.
307+
2. **Chains.** The codecs are assembled into the returned handle. Codec chains
308+
live on that handle rather than on the session, so this step writes nothing
309+
even though it can fail on a bad capsule or a duplicate id.
310+
3. **Resolve.** Every declared function is wrapped and every name is checked,
311+
and every `__datafusion_session_planner__` runs against the completed
312+
chains.
313+
4. **Commit.** The planner is bound and the functions are registered.
314+
315+
Only step 4 touches the session, and every step that can fail happens before
316+
it. This is a rule for anyone extending `with_extensions`, not only a
317+
description: a new kind of component must do its fallible work — importing a
318+
capsule, resolving a name — in step 3, so that step 4 cannot raise part-way
319+
through. There is nothing to roll back to if it does. The returned handle
320+
shares one session with the receiver, and undoing a registration is not the
321+
same as restoring what it displaced: deregistering a function that shadowed a
322+
built-in removes the built-in too.
323+
324+
(extension_bundles_collisions)=
325+
326+
### Two bundles claiming one name
327+
328+
Within a single call, two extensions declaring a function of the same kind
329+
under the same name is a `ValueError` naming both. Codec ids dispatch on
330+
decode, so a chain can hold many and pick the right one; a function registry
331+
has no such fall-through, and the second registration would silently replace
332+
the first. Names are compared per kind, so a scalar function and an aggregate
333+
may share one.
334+
335+
Shadowing a name the session *already* has is allowed and is not a collision.
336+
The registry holds every DataFusion built-in, and overriding built-ins by name
337+
is a supported thing to do — `enable_spark_functions` is built on it.
292338

293339
Like every other derivation, the returned context is a handle on the *same*
294340
session as the receiver — see {ref}`extension_sessions`. Only the Python-side

‎docs/source/extension-guide/checklist.md‎

Lines changed: 6 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -59,13 +59,15 @@ publish. Each links to the page that explains it.
5959

6060
## Bundles and planners
6161

62-
- [ ] **You ship a bundle, not loose pieces**, if you have codecs or a planner.
62+
- [ ] **You ship a bundle, not loose pieces**, if you have codecs, functions,
63+
or a planner.
6364
→ {ref}`extension_bundles`
6465
- [ ] **Your bundle is configuration-only.** Fresh components on every call,
6566
no cached bound components, no retaining the context passed in, no
66-
registering anything on it — a factory that mutates the context is not
67-
rolled back if a later factory raises.
68-
→ {ref}`extension_bundles`
67+
registering anything on it — declare what you contribute instead, so the
68+
host can validate it before anything is written and install it after
69+
every codec is in place.
70+
→ {ref}`extension_bundles_transaction`
6971
- [ ] **Your codecs are objects exposing the getter, not bare capsules.**
7072
`with_extensions` refuses a capsule, because there would be nothing to
7173
name the codec by. → {ref}`extension_bundles_codecs_are_objects`

‎docs/source/extension-guide/functions.md‎

Lines changed: 22 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -26,12 +26,12 @@ functions in pure Python — see
2626
{doc}`../user-guide/common-operations/udf-and-udfa` — and the two roads meet at
2727
the same registration methods.
2828

29-
| Hook | Contributes | Wrapped by | Registered with |
30-
| --- | --- | --- | --- |
31-
| `__datafusion_scalar_udf__` | scalar function | {py:func}`datafusion.udf` | {py:meth}`~datafusion.SessionContext.register_udf` |
32-
| `__datafusion_aggregate_udf__` | aggregate function | {py:func}`datafusion.udaf` | {py:meth}`~datafusion.SessionContext.register_udaf` |
33-
| `__datafusion_window_udf__` | window function | {py:func}`datafusion.udwf` | {py:meth}`~datafusion.SessionContext.register_udwf` |
34-
| `__datafusion_table_function__` | function returning a table | {py:func}`datafusion.udtf` | {py:meth}`~datafusion.SessionContext.register_udtf` |
29+
| Hook | Contributes | Wrapped by | Registered with | Declared in a bundle as |
30+
| --- | --- | --- | --- | --- |
31+
| `__datafusion_scalar_udf__` | scalar function | {py:func}`datafusion.udf` | {py:meth}`~datafusion.SessionContext.register_udf` | `udfs` |
32+
| `__datafusion_aggregate_udf__` | aggregate function | {py:func}`datafusion.udaf` | {py:meth}`~datafusion.SessionContext.register_udaf` | `udafs` |
33+
| `__datafusion_window_udf__` | window function | {py:func}`datafusion.udwf` | {py:meth}`~datafusion.SessionContext.register_udwf` | `udwfs` |
34+
| `__datafusion_table_function__` | function returning a table | {py:func}`datafusion.udtf` | {py:meth}`~datafusion.SessionContext.register_udtf` | — |
3535

3636
All four are implemented in [`datafusion-ffi-example`], one per file.
3737

@@ -67,6 +67,22 @@ from datafusion import udf
6767
ctx.register_udf(udf(my_library.MyScalarUDF()))
6868
```
6969

70+
If your library ships more than a function or two, declare them on a bundle
71+
instead and let one call install everything:
72+
73+
```python
74+
class MyLibraryExtension:
75+
def __datafusion_session_components__(self, ctx):
76+
return SessionExtensionComponents(udfs=(my_library.MyScalarUDF(),))
77+
78+
79+
ctx = SessionContext().with_extensions(MyLibraryExtension())
80+
```
81+
82+
Either the raw exportable or an already-wrapped
83+
{py:class}`~datafusion.user_defined.ScalarUDF` is accepted; the name comes off
84+
the capsule either way. See {ref}`extension_bundles`.
85+
7086
## Table functions
7187

7288
A table function takes literal `Expr` arguments and returns a table provider,

‎docs/source/user-guide/extensions.md‎

Lines changed: 9 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -36,9 +36,8 @@ this repository under
3636

3737
Which one you have determines how much setup you do.
3838

39-
**Tables and functions register directly.** If the library gives you a table
40-
or a function, register it the same way you would register a CSV file. No
41-
extra setup:
39+
**Tables register directly.** If the library gives you a table, register it
40+
the same way you would register a CSV file. No extra setup:
4241

4342
```python
4443
from datafusion import SessionContext
@@ -68,6 +67,13 @@ ctx.sql("SELECT count(*) FROM events").show()
6867
everything else with the context you called it on, so tables you registered
6968
before the call are still there.
7069

70+
**Functions can arrive either way.** A single function is registered directly
71+
with {py:func}`~datafusion.udf` and
72+
{py:meth}`~datafusion.SessionContext.register_udf`. A library shipping a set of
73+
them usually packages them in the same `Extension` object instead, so
74+
`with_extensions` installs them along with everything else it provides. Follow
75+
whichever the library documents.
76+
7177
## Using more than one library
7278

7379
Pass them all to a single call:

‎examples/datafusion-ffi-example/README.md‎

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -29,6 +29,10 @@ The example intentionally uses separate `cdylib` crates for these roles:
2929

3030
Separate shared libraries guarantee distinct DataFusion library markers. This catches type-identity mistakes that a planner and provider compiled into one shared library would hide.
3131

32+
## Installing the functions as a bundle
33+
34+
`MyFunctionExtension` implements `__datafusion_session_components__` and declares this crate's scalar, aggregate, and window functions, so a caller installs all three with one `SessionContext.with_extensions(MyFunctionExtension())` rather than wrapping and registering each in turn. It contributes no codecs and no planner, which is the shape a function-only library takes. `python/tests/_test_session_extension.py` covers it, including that a failure after the hook registers nothing.
35+
3236
## Codec behavior
3337

3438
`MyLogicalExtensionCodec` serializes this example's in-memory table providers, and `MyPhysicalExtensionCodec` serializes provider-owned memory scans and opaque FFI wrappers around them. Both use documented, process-local, one-shot token registries. The registries make ownership and callback routing visible without pretending to be a portable format. They assume trusted in-process payloads and consume each token during decoding. A production provider should instead encode durable metadata from which its provider and plans can be reconstructed.
Lines changed: 128 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,128 @@
1+
# Licensed to the Apache Software Foundation (ASF) under one
2+
# or more contributor license agreements. See the NOTICE file
3+
# distributed with this work for additional information
4+
# regarding copyright ownership. The ASF licenses this file
5+
# to you under the Apache License, Version 2.0 (the
6+
# "License"); you may not use this file except in compliance
7+
# with the License. You may obtain a copy of the License at
8+
#
9+
# http://www.apache.org/licenses/LICENSE-2.0
10+
#
11+
# Unless required by applicable law or agreed to in writing,
12+
# software distributed under the License is distributed on an
13+
# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
14+
# KIND, either express or implied. See the License for the
15+
# specific language governing permissions and limitations
16+
# under the License.
17+
18+
"""What a function library gets for shipping a bundle instead of a recipe."""
19+
20+
from __future__ import annotations
21+
22+
import pyarrow as pa
23+
import pytest
24+
from datafusion import SessionContext, SessionExtensionComponents
25+
from datafusion_ffi_example import MyFunctionExtension
26+
27+
28+
def _session():
29+
"""A session with the library installed in one call.
30+
31+
The comparison this file exists to make: without a bundle this is three
32+
``register_*`` calls the caller has to know about, one per function.
33+
"""
34+
ctx = SessionContext().with_extensions(MyFunctionExtension())
35+
batch = pa.RecordBatch.from_arrays([pa.array([1, 2, 3, None])], names=["a"])
36+
ctx.register_record_batches("test_table", [[batch]])
37+
return ctx
38+
39+
40+
def test_one_call_installs_every_function():
41+
"""All three kinds arrive across the FFI boundary from one hook."""
42+
ctx = _session()
43+
44+
scalar = ctx.sql("select my_custom_is_null(a) from test_table").collect()
45+
assert [r.column(0) for r in scalar] == [
46+
pa.array([False, False, False, True], type=pa.bool_())
47+
]
48+
49+
aggregate = ctx.sql("select my_custom_sum(a) from test_table").collect()
50+
assert aggregate[0].column(0)[0].as_py() == 6
51+
52+
window = ctx.sql(
53+
"select my_custom_rank() over (order by a) from test_table"
54+
).collect()
55+
assert window[0].num_rows == 4
56+
57+
58+
def test_the_names_come_from_the_capsules():
59+
"""Not from anything the bundle or the host said.
60+
61+
The wrappers are built by the host during resolution, so a name it invented
62+
would be the one a query had to use. These are the names the Rust
63+
``ScalarUDFImpl`` and friends report.
64+
"""
65+
ctx = _session()
66+
67+
assert ctx.udf("my_custom_is_null").name == "my_custom_is_null"
68+
assert ctx.udaf("my_custom_sum").name == "my_custom_sum"
69+
assert ctx.udwf("my_custom_rank").name == "my_custom_rank"
70+
71+
72+
def test_the_bundle_is_reusable_across_sessions():
73+
"""One bundle object, two sessions: components are built per install."""
74+
extension = MyFunctionExtension()
75+
first = SessionContext().with_extensions(extension)
76+
second = SessionContext().with_extensions(extension)
77+
78+
assert first.udf("my_custom_is_null").name == "my_custom_is_null"
79+
assert second.udf("my_custom_is_null").name == "my_custom_is_null"
80+
assert first.session_id() != second.session_id()
81+
82+
83+
def test_installing_the_library_twice_is_refused():
84+
"""The collision rule holds for functions arriving over FFI.
85+
86+
Two instances of one library is the shape this actually takes in the wild —
87+
an application assembling its extension list from a plugin registry that
88+
lists the same package twice.
89+
"""
90+
ctx = SessionContext()
91+
92+
with pytest.raises(ValueError, match=r"scalar function named 'my_custom_is_null'"):
93+
ctx.with_extensions(MyFunctionExtension(), MyFunctionExtension())
94+
95+
with pytest.raises(KeyError):
96+
ctx.udf("my_custom_is_null")
97+
98+
99+
def test_a_failure_after_the_hook_registers_nothing():
100+
"""The transaction covers functions imported across the FFI boundary too."""
101+
ctx = SessionContext()
102+
103+
class BoomPlanner:
104+
def __datafusion_session_planner__(self, ctx, fallback) -> None:
105+
msg = "boom"
106+
raise RuntimeError(msg)
107+
108+
with pytest.raises(RuntimeError, match="boom"):
109+
ctx.with_extensions(MyFunctionExtension(), BoomPlanner())
110+
111+
with pytest.raises(KeyError):
112+
ctx.udf("my_custom_is_null")
113+
114+
115+
def test_the_hook_returns_the_components_type():
116+
"""The bundle builds a real dataclass, not a duck-typed stand-in.
117+
118+
``with_extensions`` rejects anything else, so a Rust bundle that imported
119+
the wrong name would fail at install rather than silently contribute
120+
nothing.
121+
"""
122+
components = MyFunctionExtension().__datafusion_session_components__(
123+
SessionContext()
124+
)
125+
126+
assert isinstance(components, SessionExtensionComponents)
127+
assert len(components.udfs) == 1
128+
assert components.logical_extension_codecs == ()

‎examples/datafusion-ffi-example/src/aggregate_udf.rs‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -40,7 +40,7 @@ pub(crate) struct MySumUDF {
4040
#[pymethods]
4141
impl MySumUDF {
4242
#[new]
43-
fn new() -> PyResult<Self> {
43+
pub(crate) fn new() -> PyResult<Self> {
4444
Ok(Self {
4545
inner: Arc::new(Sum::new()),
4646
})

0 commit comments

Comments
 (0)