Skip to content

Commit a63583c

Browse files
timsaucerclaude
andcommitted
Add dfx_udfs: the library that cannot be installed as a bundle
Second of three libraries for #1719, and the one carrying the mixed-workflow case: it exposes no `__datafusion_session_components__`, so callers register its three functions and install its two codecs by hand. That is not an artificial handicap. `SessionExtensionComponents` carries codec fields only, so a function library has nowhere to put its functions — the rename in 5a1bfeb noted that UDF and provider fields will join later. Until they do, this is what a function library actually looks like, and the example should show what that costs rather than pretend every dependency has caught up. A test asserts the shape rather than describing it: `with_extensions` rejects this object, naming the hook it lacks. The functions are `dfx_net_revenue` (the TPC-H revenue expression), `dfx_weighted_avg`, and `dfx_revenue_rank`. The aggregate is written out rather than delegating to a built-in because its state is the point: two running sums, which is what lets DataFusion compute a partial aggregate per partition and merge the results. An aggregate that could only be evaluated over its whole input at once would give a different answer once split, which is exactly what a worker does to it. A test pins that by checking the plan really is `mode=Partial` and the answer is still right. Both codecs are name-only — `try_encode_*` writes nothing and `try_decode_*` rebuilds from `name`. Three worker tests, each a separate interpreter, pin what that buys, and they disagree with each other in the useful way: codec installed, nothing registered -> works, decode_calls == 1 functions registered, no codec -> works, decode_calls == 0 neither -> fails, naming dfx_net_revenue So installing the codec is an *alternative* to registering the functions, not an addition to it. The middle case is the trap worth knowing: on the driver, where the functions are registered, the registry is tried first and the codec is never consulted — so a codec that was broken or missing looks fine right up until a worker needs it. `try_decode_*` checks the name before the buffer, in that order. An empty encoding leaves `fun_definition` unset and carries no codec id, so it is the one path where a payload is offered to every installed codec in turn; a codec that trusted `buf` first would answer for names it does not own. There are two codecs because there are two plan layers and a library cannot know which one its callers will serialize — an engine shipping physical plans exercises only the physical one. Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
1 parent 29502bf commit a63583c

11 files changed

Lines changed: 1304 additions & 0 deletions

File tree

‎Cargo.lock‎

Lines changed: 18 additions & 0 deletions
Some generated files are not rendered by default. Learn more about customizing how changed files appear on GitHub.

‎Cargo.toml‎

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -33,6 +33,7 @@ members = [
3333
"examples/datafusion-ffi-example",
3434
"examples/datafusion-ffi-query-planner-example",
3535
"examples/distributed/storage-library",
36+
"examples/distributed/udf-library",
3637
]
3738
resolver = "3"
3839

Lines changed: 51 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,51 @@
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+
[package]
19+
name = "dfx-udfs"
20+
version.workspace = true
21+
edition.workspace = true
22+
rust-version.workspace = true
23+
license.workspace = true
24+
description = "Example extension library: user defined functions plus the codecs that make them portable"
25+
homepage.workspace = true
26+
repository.workspace = true
27+
publish = false
28+
29+
[dependencies]
30+
arrow = { workspace = true }
31+
arrow-schema = { workspace = true }
32+
datafusion = { workspace = true }
33+
datafusion-common = { workspace = true, default-features = false }
34+
datafusion-expr = { workspace = true }
35+
datafusion-ffi = { workspace = true }
36+
datafusion-functions-window = { workspace = true }
37+
datafusion-proto = { workspace = true }
38+
datafusion-python-util.workspace = true
39+
pyo3 = { workspace = true, features = [
40+
"extension-module",
41+
"abi3",
42+
"abi3-py310",
43+
] }
44+
pyo3-log = { workspace = true }
45+
46+
[build-dependencies]
47+
pyo3-build-config = { workspace = true }
48+
49+
[lib]
50+
name = "dfx_udfs"
51+
crate-type = ["cdylib", "rlib"]
Lines changed: 20 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,20 @@
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+
fn main() {
19+
pyo3_build_config::add_extension_module_link_args();
20+
}
Lines changed: 32 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,32 @@
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+
[build-system]
19+
requires = ["maturin>=1.6,<2.0"]
20+
build-backend = "maturin"
21+
22+
[project]
23+
name = "dfx_udfs"
24+
requires-python = ">=3.10"
25+
classifiers = [
26+
"Programming Language :: Rust",
27+
"Programming Language :: Python :: Implementation :: CPython",
28+
]
29+
dynamic = ["version"]
30+
31+
[tool.maturin]
32+
features = ["pyo3/extension-module"]

0 commit comments

Comments
 (0)