Conversation
…nctions Close gaps found by auditing the Python API against upstream DataFusion 55.1.0. - Add the any_value aggregate. - Add array_add, array_subtract, array_scale, array_sum, array_avg, array_product, and array_first, each with its list_* alias. - Add Spark monthname, weekday, atan2, hypot, pow/power, quote, and concat_ws. concat_ws calls the UDF directly so *cols stays variadic. - Fix the aggregations guide, which listed regr_slope twice and omitted regr_sxy. Co-Authored-By: Claude Opus 5.5 <[email protected]>
…t kwargs - Add rand and substring_index as aliases of random and substr_index. - Add input_file_name and file_row_index, which report the source file and row offset during a file scan. - Add a distinct argument to bit_and, bit_or, mean, percentile_cont, quantile_cont, and string_agg. distinct goes before filter, matching sum and avg; the upgrade guide covers positional callers. - Fix mean, which passed filter into avg's distinct slot and raised a TypeError whenever filter was given. Co-Authored-By: Claude Opus 5.5 <[email protected]>
- Add getbit, dateadd, datediff, datepart, sha, ceiling, printf, char_length, and character_length as aliases of their Spark primaries. - Add substr, whose len argument is optional as in pyspark. It calls the UDF directly because upstream expr_fn::substring always takes a length. - Rename the spark.last_day parameter from col to date to match pyspark, with an upgrade-guide note for keyword callers. Co-Authored-By: Claude Opus 5.5 <[email protected]>
…tring, substr - btrim, ltrim, rtrim, and trim take an optional characters argument naming the set to strip. - array_to_string and its aliases take an optional null_string that is written in place of NULL elements. - substr takes an optional length. Co-Authored-By: Claude Opus 5.5 <[email protected]>
Replace NaN in floating-point columns, optionally limited to a subset. Mirrors fill_null and wraps upstream DataFrame::fill_nan. Co-Authored-By: Claude Opus 5.5 <[email protected]>
IGNORE_NULLS skips null values when counting shift_offset rows, matching SQL LEAD/LAG ... IGNORE NULLS. The argument is appended after order_by, so existing calls are unaffected. Co-Authored-By: Claude Opus 5.5 <[email protected]>
5f33af0 to
7cf8fb5
Compare
…xplain Expose the remaining upstream ExplainOption fields as keywords on DataFrame.explain. Each defaults to None, which falls back to the matching datafusion.explain.* session setting, so existing calls are unaffected. New ExplainAnalyzeLevel and ExplainMetricCategory enums sit beside ExplainFormat. Co-Authored-By: Claude Opus 5.5 <[email protected]>
Upstream ExprFunctionExt methods on an Expr start from an empty builder, so build() resets every option not set again. The Python function wrappers already apply their keyword options, so calls such as string_agg(..., order_by=...).distinct().build() silently dropped the ordering, and .filter() dropped order_by, and so on. Seed the builder from the expression's existing params instead. A window frame equal to the default for its order_by is left unset so build() derives it again from the final order_by. Co-Authored-By: Claude Opus 5.5 <[email protected]>
AggregateUDF already accepted the capsule returned by __datafusion_aggregate_udf__ as well as an object exposing it (apache#1277). Extend the same to ScalarUDF and WindowUDF, with matching overloads on udf and udwf, so the three UDF kinds import the same way. Add FFI example tests for all three, including the previously untested AggregateUDF path. Co-Authored-By: Claude Opus 5.5 <[email protected]>
The ABC is defined in datafusion.catalog beside CatalogProvider and SchemaProvider, which are in its __all__, but was only listed in the package root's __all__. Co-Authored-By: Claude Opus 5.5 <[email protected]>
str(df) is the Pythonic way to get a string and already uses __repr__ and the configurable formatter, so check-upstream should not flag the upstream to_string as a gap. Co-Authored-By: Claude Opus 5.5 <[email protected]>
7cf8fb5 to
5d43576
Compare
user_defined.py imported CapsuleType from _typeshed, which does not define it. Pyright reports `"CapsuleType" is unknown import symbol`, so _PyCapsule resolved to Unknown and every overload and parameter typed with it accepted any argument. This affected the udaf/from_pycapsule hints added in apache#1277 as well as the new udf/udwf ones. Import it from types on Python 3.13+ and from typing_extensions (already a dependency below 3.13) otherwise. Checked with pyright at --pythonversion 3.10 and 3.13: udf/udaf/udwf and each from_pycapsule accept a CapsuleType and return the right wrapper, and WindowUDF.from_pycapsule(1) is now rejected where it previously passed. Co-Authored-By: Claude Opus 5.5 <[email protected]>
A built window function always stores a concrete frame, so the builder could not tell a frame the user chose from the default. A frame equal to the no-order_by default was treated as unset and re-derived as the running frame once order_by was chained, silently changing results. Record on the Python Expr whether over() or window_frame() set the frame explicitly, and pass that to builder_from_expr so the frame is kept. Co-Authored-By: Claude Opus 5.5 <[email protected]>
Defaulting to NullTreatment.RESPECT_NULLS passed Some(RespectNulls) to Rust, which adds "RESPECT NULLS" to the generated column name and breaks code that refers to an un-aliased lead/lag output by name. Default to None instead; respecting nulls is already the behavior when unset. Co-Authored-By: Claude Opus 5.5 <[email protected]>
Calling over() on an expression that was already a window function rebuilt it from an empty builder, silently dropping its order_by, null_treatment, and explicit window frame. For example, lead(v, order_by="i", null_treatment=IGNORE_NULLS).over(Window(partition_by=[g])) lost both the ordering and IGNORE NULLS. Start from builder_from_expr instead, and apply only the options the Window sets. Pass the Python-side explicit-frame flag through so a frame equal to the default is not re-derived when over() adds an order_by. Co-Authored-By: Claude Opus 5.5 <[email protected]>
range() required all three of start, stop, and step as Expr, although upstream also accepts range(stop) and range(start, stop). Make stop and step optional in the bindings for both range and gen_series, so a single argument is the upper bound starting at 0, like Python's built-in range. Accept plain ints for start, stop, and step, coerced to literals, so callers no longer need lit(). Passing step without stop raises ValueError. Co-Authored-By: Claude Opus 5.5 <[email protected]>
There was a problem hiding this comment.
Copilot review overview
🟡 Changes recommended
Public explain enums need consistent exports, and identified documentation and distinct-option coverage gaps remain unresolved.
Get a fresh assessment by requesting another Copilot review.
Review effort: Balanced
Findings: 1
Open (3)
What changed in this PR
Expands Python coverage for DataFusion 55.1.0 APIs and fixes option preservation across expression builders.
Changes:
- Adds array, aggregate, metadata, Spark, range, and DataFrame APIs.
- Preserves aggregate/window builder options and supports bare FFI capsules.
- Adds documentation, migration guidance, and Python integration tests.
| File | Description |
|---|---|
.ai/skills/check-upstream/SKILL.md |
Records to_string audit guidance. |
crates/core/src/dataframe.rs |
Binds explain options and fill_nan. |
crates/core/src/expr.rs |
Preserves expression-builder options. |
crates/core/src/functions.rs |
Binds new functions and arguments. |
crates/core/src/spark_functions.rs |
Adds Spark function bindings. |
crates/core/src/udf.rs |
Imports bare scalar UDF capsules. |
crates/core/src/udwf.rs |
Imports bare window UDF capsules. |
docs/source/user-guide/common-operations/aggregations.md |
Updates aggregate catalog. |
docs/source/user-guide/upgrade-guides.md |
Documents breaking signature changes. |
examples/datafusion-ffi-example/python/tests/_test_aggregate_udf.py |
Tests bare aggregate capsules. |
examples/datafusion-ffi-example/python/tests/_test_scalar_udf.py |
Tests bare scalar capsules. |
examples/datafusion-ffi-example/python/tests/_test_window_udf.py |
Tests bare window capsules. |
python/datafusion/catalog.py |
Exports TableProviderFactory. |
python/datafusion/dataframe.py |
Exposes explain options and fill_nan. |
python/datafusion/expr.py |
Tracks explicit window frames. |
python/datafusion/functions/__init__.py |
Adds functions and optional arguments. |
python/datafusion/functions/spark.py |
Adds Spark functions and aliases. |
python/datafusion/user_defined.py |
Accepts bare UDF capsules. |
python/tests/test_aggregation.py |
Covers aggregate additions. |
python/tests/test_dataframe.py |
Covers DataFrame and window options. |
python/tests/test_expr.py |
Covers builder option preservation. |
python/tests/test_functions.py |
Covers scalar and array additions. |
python/tests/test_lambda.py |
Covers array_first. |
python/tests/test_spark_functions.py |
Covers Spark additions and aliases. |
💡 Add a code-review agent skill or configure MCP servers for context-aware, tailored reviews. Learn more in the docs.
andygrove
left a comment
There was a problem hiding this comment.
Thanks for this, it closes a lot of gaps. Each inline comment has a repro, run against both this PR's head (96f990b) and its merge base (90f8709).
The biggest theme is the new keep-existing-options behavior when chaining builder methods or .over(). It makes some things work that didn't before (filter/distinct on an aggregate window now build correctly), but:
- options that don't apply to the function's kind are now silently dropped instead of raising (
builder_from_expr), - the explicit-frame flag lives only on the Python object, so it's lost through copy/pickle/
from_bytes/SQL, - the result and column-name changes aren't in the upgrade guide yet.
A few comments are about bugs that predate this PR but sit in code it touches (percentile_cont's sort, .over() on aggregates / #1764, the keyword-only udf decorator). Those are fine to split out.
One more that can't go inline because the file isn't in the diff: skills/datafusion_python/SKILL.md (lines 747-760) still says F.substr takes only two arguments, that a third raises TypeError, and labels F.substr(col("c_phone"), lit(1), lit(2)) as WRONG. With the new length argument that call works (returns '25'), so the skill now steers agents away from a valid form.
| }; | ||
|
|
||
| let cols = cols.iter().map(String::as_str).collect::<Vec<_>>(); | ||
| let df = self.df.as_ref().fill_nan(&scalar_value.0, &cols)?; |
There was a problem hiding this comment.
Upstream DataFrame::fill_nan (via fill_columns) rebuilds every column with col(field.name()), which parses the name as an identifier (lowercasing it, splitting on .). So fill_nan fails on any DataFrame with an uppercase or dotted column name, even when that column isn't in subset:
from datafusion import SessionContext
nan = float("nan")
df = SessionContext().from_pydict({"Name": ["a", "b"], "v": [nan, 1.0]})
df.fill_nan(0.0, subset=["v"])
# Schema error: No field named name. Did you mean '..."Name"'?fill_null has the same upstream bug, so this is probably worth an upstream issue (ident(field.name()) there would fix both). Until then the wrapper could build the projection itself from unparsed column references plus nanvl.
There was a problem hiding this comment.
Confirmed in pure Rust on 55.1.0, for both Name and a.b column names, and filed upstream as apache/datafusion#25829 with a suggested fix (Expr::Column(Column::from((qualifier, field))), which also keeps the qualifier). I've left the wrapper as-is so fill_nan and fill_null keep behaving the same way, and we'll pick up the upstream fix when it lands.
The explicit-frame flag added in 6345b17 lived only on the Python Expr, so copy, pickle, from_bytes and SQL all lost it. The same builder chain then gave different results for identical expressions, including between a driver and the workers it pickles expressions to. Drop the flag. A frame equal to the default for the current order-by is treated as unset when chaining, which matches main. Document the rule and how to keep such a frame (set it after order_by, or in the same Window) in the window functions guide. Co-Authored-By: Claude Opus 5.5 (1M context) <[email protected]>
…ning build() derives RANGE UNBOUNDED PRECEDING .. CURRENT ROW from an empty order_by list and the whole-partition frame from an absent one, but both are stored as an empty order_by. Only the whole-partition frame was recognised as a default, so a frame derived from order_by=[] looked explicit and survived a later order_by, summing tied rows together. Treat either frame as the default when the stored order_by is empty. Co-Authored-By: Claude Opus 5.5 (1M context) <[email protected]>
Seeding the builder from the existing expression always gives it a function kind, so upstream's per-method kind checks no longer applied. partition_by and window_frame on a plain aggregate, and distinct on a non-aggregate window function, built successfully and silently dropped the option instead of raising at build(). Fall back to upstream's ExprFunctionExt method for those cases so build() raises again. distinct on an aggregate run as a window function still keeps the existing options. Co-Authored-By: Claude Opus 5.5 (1M context) <[email protected]>
Expr.over() on an aggregate drops its order_by, null_treatment, filter, and distinct options (apache#1764). The new distinct arguments on mean, percentile_cont, quantile_cont, and string_agg run straight into it. Document the limitation and the workaround (pass order_by and null_treatment in the Window, chain filter() or distinct() after over()) in the window functions guide, and point to it from Expr.over. Co-Authored-By: Claude Opus 5.5 (1M context) <[email protected]>
Add an upgrade-guide section for the behavior change from keeping existing options when chaining builder methods or over(): results can change (lead keeps its partition), chains whose options used to be dropped can now raise (array_agg with distinct and order_by), and first_value, last_value, and nth_value gain RESPECT NULLS in their generated names. Rewrite the Expr.over docstring, which described only aggregates, to cover window functions and point to the frame-chaining rule. Co-Authored-By: Claude Opus 5.5 (1M context) <[email protected]>
andygrove
left a comment
There was a problem hiding this comment.
Thanks for working through all of these. Each inline comment has a repro, run against both this PR's head (db5b7cc) and its merge base (90f8709).
Most of what's left is in last round's fixes, where the fix covers the case I gave but not a close neighbor:
percentile_contand friends keep the direction as aggregates, but.over()still drops it, and thewindows.mdworkaround can't bring it back (1b01845),- the empty-
order_byframe (b4cad2d) is kept after copy/pickle/from_bytes, because decoding addsORDER BY UInt64(1), - the empty-
subsetguard (8454a09) raises on numpy arrays andpandas.Indexsubsets, which worked on main, - the kind check (92ace16) only covers the first builder call, and
filter()has none, explain(607dde4) still ignoresshow_statistics=Falsewithanalyze, andanalyze_level/analyze_categorieswithout it,- the #1764 workaround in
windows.md(8e21c8a) gives the non-distinct result for Python UDAFs, string_agg's bool check (0ac494e) rejectsnp.True_,udf(func=...)with the callable passed by keyword still fails (f8d7b70).
The rest are new: range(start=5), a few docstrings that promise more than the code does, and the CapsuleType import in context.py/extensions.py. The UDAF distinct, udf(func=...), and count_star ones predate this PR and are fine to split out.
| @@ -678,8 +705,8 @@ impl PyExpr { | |||
| null_treatment, | |||
There was a problem hiding this comment.
1b01845 keeps the direction for the aggregate, but this arm still builds the window from agg_fn.func and agg_fn.params.args only (lines 695-698), so the WITHIN GROUP ordering is dropped again and a descending percentile used as a window silently returns the ascending result. The windows.md workaround (pass order_by in the Window) can't bring it back, since a window never hands its ORDER BY to the accumulator:
from datafusion import SessionContext, col, functions as f
from datafusion.expr import Window, WindowFrame
ctx = SessionContext()
df = ctx.from_pydict({"a": [1.0, 2.0, 3.0, 4.0, 5.0]}, name="t")
desc = col("a").sort(ascending=False)
whole = WindowFrame("rows", None, None)
print(df.aggregate([], [f.percentile_cont(desc, 0.25).alias("p")]).to_pydict())
# {'p': [4.0]}
print(df.select(f.percentile_cont(desc, 0.25).over(Window()).alias("p")).to_pydict())
# {'p': [2.0, 2.0, 2.0, 2.0, 2.0]}
print(df.select(
f.percentile_cont(desc, 0.25).over(Window(order_by=[desc], window_frame=whole)).alias("p")
).to_pydict())
# {'p': [2.0, 2.0, 2.0, 2.0, 2.0]}
ctx.sql("SELECT percentile_cont(0.25) WITHIN GROUP (ORDER BY a DESC) OVER () FROM t")
# Error during planning: OVER and WITHIN GROUP clause cannot be used together. ...quantile_cont gives the same, and approx_percentile_cont gives 4.25 as an aggregate vs 1.75 over a window. SQL also rejects a plain aggregate ORDER BY with OVER ("Aggregate ORDER BY is not implemented for window functions"), so raising here when agg_fn.params.order_by is non-empty would match SQL and turn this, and the order_by part of #1764, into an error instead of a wrong answer. The upgrade guide's "They now match WITHIN GROUP (ORDER BY ... DESC)" could also say it's for aggregates only.
There was a problem hiding this comment.
Fixed in 2fb7cce: fix: keep aggregate options in over() and reject an aggregate order_by. over() now keeps filter, distinct, and null_treatment, and an aggregate order_by raises, matching SQL. A WITHIN GROUP function accepts an ascending sort_expression and raises on a descending one, so the descending percentile case is now an error instead of a wrong answer. The upgrade guide says the percentile change applies to aggregates, and 4d66257 (docs: show how to move an aggregate's order_by into the Window) adds a before/after for the #1764 case.
| // left unset, so it is derived again from the final order-by. An | ||
| // absent and an empty order-by derive different frames but are both | ||
| // stored as empty, so either frame counts as the default here. | ||
| let is_default_frame = if has_order_by { |
There was a problem hiding this comment.
b4cad2d fixed the chain I gave, but this check keys off params.order_by.is_empty(), and decoding (copy, pickle, from_bytes) runs upstream's regularize_order_bys, which adds ORDER BY UInt64(1) to a RANGE frame that has no order-by. After a round trip the frame no longer looks like the default, so it's kept, and the copy chains differently from the original:
import copy
import pickle
from datafusion import Expr, SessionContext, col, functions as f
from datafusion.expr import Window
ctx = SessionContext()
df = ctx.from_pydict({"i": [1, 1, 2, 3], "v": [1, 2, 3, 4]})
def run(e):
r = e.over(Window(order_by="i")).alias("r")
return df.select(col("v"), r).sort(col("v")).collect_column("r").to_pylist()
base = f.sum(col("v")).over(Window(order_by=[]))
print(run(base)) # [1, 3, 6, 10]
print(run(copy.copy(base))) # [3, 3, 6, 10]
print(run(pickle.loads(pickle.dumps(base)))) # [3, 3, 6, 10]
print(run(Expr.from_bytes(base.to_bytes()))) # [3, 3, 6, 10]
print(copy.copy(base))
# Expr(sum(v) ORDER BY [UInt64(1) ASC NULLS LAST] RANGE BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW)On main all four give [1, 3, 6, 10]. windows.md promises this case ("a copy, a pickled expression sent to a worker, or one parsed from SQL all chain the same way"), but test_window_builder_default_frame_same_after_copy_and_pickle only covers WindowFrame("rows", None, None). The check probably needs to treat an order_by holding only that literal sort key the same as an empty one.
There was a problem hiding this comment.
Fixed in 99f506d: fix: treat the sort key added when decoding a RANGE frame as no order_by
| - For columns where casting fails, the original column is kept unchanged | ||
| - For columns not in subset, the original column is kept unchanged | ||
| """ | ||
| if subset is not None and not subset: |
There was a problem hiding this comment.
not subset takes the truth value of subset, so column lists from numpy or pandas now raise here (and in fill_nan below). The Rust binding accepts any sequence of strings, so these worked on main:
import numpy as np
import pandas as pd
from datafusion import SessionContext
df = SessionContext().from_pydict({"a": [1, None, 3], "b": [None, 5, 6], "c": [None, None, 9]})
cols = pd.DataFrame(columns=["a", "b"]).columns
print(df.fill_null(0, subset=cols).to_pydict())
# main: {'a': [1, 0, 3], 'b': [0, 5, 6], 'c': [None, None, 9]}
# PR: ValueError: The truth value of a Index is ambiguous. ...
print(df.fill_null(0, subset=np.array(["a", "b"])).to_pydict())
# main: {'a': [1, 0, 3], 'b': [0, 5, 6], 'c': [None, None, 9]}
# PR: ValueError: The truth value of an array with more than one element is ambiguous. ...if subset is not None and len(subset) == 0: keeps the new empty-subset behavior without this.
There was a problem hiding this comment.
Fixed in e811758: fix: accept numpy and pandas column lists as fill_null and fill_nan subsets
| .into() | ||
| } | ||
|
|
||
| pub fn filter(&self, filter: PyExpr) -> PyExprFuncBuilder { |
There was a problem hiding this comment.
92ace16 restores the error for the first call, but the check only runs there. filter() has no check, and the chained PyExprFuncBuilder methods have none either, so the same option later in a chain is still dropped. The PR description says "A builder option that does not apply to the function's kind still raises":
from datafusion import SessionContext, col, lit, functions as f
ctx = SessionContext()
df = ctx.from_pydict({"g": [1, 1, 2, 2], "v": [1, 2, 3, 4]})
f.lead(col("v"), 1, order_by="v").distinct().build()
# ExprFunctionExt can only be used with Expr::AggregateFunction or Expr::WindowFunction
e = f.lead(col("v"), 1, order_by="v").partition_by(col("g")).distinct().build()
print(e) # Expr(lead(DISTINCT v, Int64(1), NULL) PARTITION BY [g] ORDER BY [v ASC NULLS FIRST] ...)
print(df.select(col("v"), e.alias("r")).sort(col("v")).collect_column("r").to_pylist())
# [2, None, 4, None], DISTINCT ignored
e = f.sum(col("v")).filter(col("v") > lit(1)).partition_by(col("g")).build()
print(df.aggregate([], [e.alias("s")]).to_pydict())
# {'s': [9]}, partition_by ignored
e = f.row_number(order_by=[col("v")]).filter(col("v") > lit(2)).build()
# main: raises here; PR: builds
df.select(e.alias("r"))
# Error during planning: FILTER clause can only be used with aggregate window functions. ...The two chained cases drop the option on main too, so only filter() is a regression, and it still errors, just later. Recording the function kind on PyExprFuncBuilder and checking there would cover every entry point with one rule.
There was a problem hiding this comment.
Fixed in 7dd2658: fix: check the function kind on every builder call, not just the first. The function kind is recorded on PyExprFuncBuilder and checked on every call, including filter().
| If using the builder functions described in ref:`_aggregation` this function ignores | ||
| the options ``order_by``, ``null_treatment``, and ``distinct``. | ||
| If using the builder functions described in :ref:`aggregation` this function ignores | ||
| the option ``null_treatment``, and ``order_by`` replaces the ordering of |
There was a problem hiding this comment.
Only the direction of the first order_by key is used. The percentile is still computed over sort_expression, while the output name shows the new column, so order_by doesn't replace the ordering the way this reads (same text at 5447 and 5497):
from datafusion import SessionContext, col, functions as f
ctx = SessionContext()
df = ctx.from_pydict({"a": [1.0, 2.0, 3.0, 4.0, 5.0], "b": [10.0, 20.0, 30.0, 40.0, 50.0]}, name="t")
e = f.percentile_cont(col("a"), 0.25).order_by(col("b")).build()
print(df.aggregate([], [e]).to_pydict())
# {'percentile_cont(Float64(0.25)) WITHIN GROUP [t.b ASC NULLS FIRST]': [2.0]}
print(ctx.sql("SELECT percentile_cont(0.25) WITHIN GROUP (ORDER BY b) FROM t").to_pydict())
# {'percentile_cont(Float64(0.25)) WITHIN GROUP [t.b ASC NULLS LAST]': [20.0]}The approx variants give 1.75 vs 17.5. The behavior isn't new, but the docstring now invites it. Saying that order_by only sets the direction and should use the same expression would be accurate.
There was a problem hiding this comment.
Fixed in daac7f2: docs: say a chained order_by sets only the percentile's sort direction
| >>> result.collect_column("s")[0].as_py() | ||
| 'x,y' | ||
| """ | ||
| if not isinstance(distinct, bool): |
There was a problem hiding this comment.
This also rejects numpy.bool_, which every other aggregate's distinct accepts through PyO3's Option<bool>, and the message then says "got bool":
import numpy as np
from datafusion import SessionContext, col, functions as f
df = SessionContext().from_pydict({"s": ["x", "y", "x"]})
print(df.aggregate([], [f.count(col("s"), distinct=np.True_).alias("n")]).to_pydict())
# {'n': [2]}
f.string_agg(col("s"), ",", distinct=np.True_)
# TypeError: distinct must be a bool, got bool; pass filter and order_by by keywordA shifted positional call can only put None or an Expr here, so rejecting just those two would keep the fix from 0ac494e and leave everything else to PyO3.
There was a problem hiding this comment.
Fixed in 47359df: fix: accept numpy.bool_ as string_agg's distinct
| ): | ||
| return ScalarUDF.from_pycapsule(args[0]) | ||
|
|
||
| if args and callable(args[0]): |
There was a problem hiding this comment.
Pre-existing, and f8d7b70 fixed the keyword-only decorator form, but the function form with the callable passed by keyword (declared by the second @overload) still falls through to the decorator factory. It now fails with a TypeError naming an internal helper instead of the old IndexError:
import pyarrow as pa
import pyarrow.compute as pc
from datafusion import udf
def double(x: pa.Array) -> pa.Array:
return pc.multiply(x, 2)
udf(func=double, input_fields=[pa.int64()], return_field=pa.int64(), volatility="immutable")
# main: IndexError: tuple index out of range
# PR: TypeError: ScalarUDF.udf.<locals>._decorator() got an unexpected keyword argument 'func'udaf(accum=...) and udwf(func=...) fail the same way. Taking func (accum for udaf) from kwargs when args is empty would cover it.
There was a problem hiding this comment.
Fixed in 1475123: fix: accept the callable by keyword in udf, udaf, udwf, and udtf
| """ | ||
| if subset is not None and not subset: | ||
| return self | ||
| return DataFrame(self.df.fill_nan(value, subset)) |
There was a problem hiding this comment.
On apache/datafusion#25829: leaving the wrapper as-is makes sense, but right now the limitation is only in the PR description. It hits any DataFrame with such a column, even when that column isn't in subset or isn't a float, so a note in the guide with a one-line pointer from the fill_nan/fill_null docstrings would save people some confusion:
from datafusion import SessionContext
nan = float("nan")
df = SessionContext().from_pydict({"Price": [nan, 2.0], "qty": [nan, 1.0]})
df.fill_nan(0.0, subset=["qty"])
# Schema error: No field named price. Did you mean '..."Price"'?There was a problem hiding this comment.
Fixed in 354db9a: docs: note that fill_null and fill_nan fail on uppercase or dotted names
| f.bit_and(column("a"), filter=my_filter) # after | ||
| ``` | ||
|
|
||
| Passing `filter` to `mean` previously raised a `TypeError`, whether passed |
There was a problem hiding this comment.
Nit: "whether passed positionally or by keyword; it now works" reads as if both now work, but a positional filter lands in the new distinct and still raises:
from datafusion import SessionContext, col, lit, functions as f
df = SessionContext().from_pydict({"v": [1.0, 2.0, 3.0]})
f.mean(col("v"), col("v") > lit(1.0))
# TypeError: 'Expr' object is not an instance of 'bool'
print(df.aggregate([], [f.mean(col("v"), filter=col("v") > lit(1.0)).alias("m")]).to_pydict())
# {'m': [2.5]}"it now works when passed by keyword" would match.
There was a problem hiding this comment.
Fixed in 9f2b326: docs: say mean's filter now works only when passed by keyword
| import sys | ||
|
|
||
| if sys.version_info >= (3, 13): | ||
| from types import CapsuleType as _PyCapsule |
There was a problem hiding this comment.
This fixes the alias here, but context.py:92 and extensions.py:52 still import CapsuleType from _typeshed, which f6d0478's message notes doesn't define it, so the capsule half of those unions is still Unknown. With pyright 1.1.414:
python/datafusion/context.py:92:27 - error: "CapsuleType" is unknown import symbol (reportAttributeAccessIssue)
python/datafusion/extensions.py:52:27 - error: "CapsuleType" is unknown import symbol (reportAttributeAccessIssue)
from datafusion import SessionContext
from datafusion.user_defined import ScalarUDF
ctx = SessionContext()
ctx.set_query_planner(1) # pyright: no error
ctx.with_logical_extension_codec(1) # pyright: no error
ScalarUDF.from_pycapsule(1) # pyright: "Literal[1]" is not assignable to "CapsuleType"Smaller: since _PyCapsule only exists under TYPE_CHECKING, typing.get_type_hints(ScalarUDF.from_pycapsule) (and WindowUDF's) now raises NameError: name '_PyCapsule' is not defined; on main it resolved.
Expr.over() on an aggregate rebuilt the window from the function and its arguments only, dropping filter, distinct, null_treatment, and order_by (apache#1764). For percentile_cont and the other WITHIN GROUP functions the order_by carries the sort direction, so a descending percentile used as a window silently returned the ascending result, and the windows.md workaround of passing order_by in the Window could not restore it. Carry filter, distinct, and null_treatment onto the window. A window never passes an ordering to the accumulator, so raise on an aggregate order_by as SQL does, except for an all-ascending ordering on a WITHIN GROUP function, which gives the same result as a window and is dropped. Replace the workaround in windows.md with the new behavior and add an upgrade guide entry. Co-Authored-By: Claude Opus 5.5 (1M context) <[email protected]>
Decoding a window expression (copy, pickle, from_bytes, or SQL) runs upstream's regularize_order_bys, which gives a RANGE frame with no order_by the constant sort key UInt64(1). The frame then no longer looked like the default for an empty order_by, so it was kept when chaining, and a copy of sum(v).over(Window(order_by=[])) summed tied rows together after a later order_by while the original did not. Treat an order_by holding only that sort key as empty, both when deciding whether the frame is a default and when carrying the order_by into the builder. Cover every default starting window and from_bytes in the round trip test, using tied values so a kept RANGE frame shows. Co-Authored-By: Claude Opus 5.5 (1M context) <[email protected]>
Seeding the builder from the existing expression bypassed upstream's kind checks, and 92ace16 restored them only for the first distinct, partition_by, and window_frame call on an Expr. filter() on a window function then built instead of raising (a regression from main), and a later call in a chain silently dropped an option that does not apply: distinct on lead, or partition_by after filter on sum. Record the function kind on ExprFuncBuilder and check it in every method: filter and distinct need an aggregate, including one used as a window function, and partition_by and window_frame need a window function. The error is raised when the option is set and names the option and the function. A builder on anything that is not a function still raises at build(), as upstream does. Co-Authored-By: Claude Opus 5.5 (1M context) <[email protected]>
…ubsets The empty-subset guard from 8454a09 used `not subset`, which raises "truth value ... is ambiguous" for a numpy array or pandas Index. The bindings accept any sequence of str, so these worked for fill_null on main. Check len(subset) == 0 instead. Co-Authored-By: Claude Opus 5.5 (1M context) <[email protected]>
607dde4 rejected show_statistics=True with analyze, but Analyze always uses the session's statistics setting, so show_statistics=False was ignored the same way. analyze_level and analyze_categories were ignored without analyze, since a plain Explain reports no metrics. SQL's EXPLAIN rejects all three. Raise a ValueError when show_statistics is set with analyze, and when analyze_level or analyze_categories is set without it. Co-Authored-By: Claude Opus 5.5 (1M context) <[email protected]>
96f990b made a single argument the upper bound, but it bound to the parameter named start, so range(start=5) silently meant stop=5 while range(stop=5) raised. Take the arguments through *args and **kwargs, with overloads for the one- and two-or-three-argument forms, and raise a TypeError when start is passed by keyword without stop. Every call that worked on main still works. The binding also chained a missing stop out of the argument list, so start and step alone built a two-argument call that read step as stop. Raise there too, rather than rely on the Python wrapper's check. Co-Authored-By: Claude Opus 5.5 (1M context) <[email protected]>
PythonFunctionAggregateUDF::accumulator never read is_distinct, and a Python Accumulator has no way to deduplicate on its own, so a Python UDAF with DISTINCT counted every row whenever the optimizer did not first rewrite the query to group by the distinct values: as a window function (including SQL's OVER), with an alias inside aggregate(), or alongside another distinct aggregate. over() now keeps distinct, which made this easier to reach. Return a not-implemented error when is_distinct is set. A plan the optimizer rewrote arrives without the flag and still runs. Co-Authored-By: Claude Opus 5.5 (1M context) <[email protected]>
1b01845 said a chained order_by replaces the ordering of sort_expression. Upstream keeps the WITHIN GROUP column in both the arguments and the order_by, and the accumulator reads only the direction from the order_by, so order_by(col("b")) still computes over sort_expression while the output name shows b. Say that it sets only the direction and should use the same expression. Co-Authored-By: Claude Opus 5.5 (1M context) <[email protected]>
…ntile approx_percentile_cont_with_weight does not ignore distinct: upstream rejects it at execution. count_star does not ignore it either: it counts the distinct values of the constant 1, so the result is 1 for any non-empty input. Say so in both docstrings. Co-Authored-By: Claude Opus 5.5 (1M context) <[email protected]>
0ac494e rejected anything that is not a Python bool, which also caught numpy.bool_, accepted by every other aggregate's distinct through PyO3, and reported it as "got bool". A shifted positional call can only put the old filter value, None or an Expr, into distinct, so reject just those and leave the rest to PyO3. Co-Authored-By: Claude Opus 5.5 (1M context) <[email protected]>
The overloads declare the function form with a named first parameter (func, or accum for udaf), but dispatch looked only at args[0], so udf(func=double, ...) fell through to the decorator factory and raised a TypeError naming an internal helper (IndexError on main). Move that parameter from kwargs into args when no positional argument is given. Co-Authored-By: Claude Opus 5.5 (1M context) <[email protected]>
Upstream rebuilds every column with col(field.name()), which parses the name as an identifier, so both methods fail on a DataFrame with any column whose name has an uppercase letter or a dot, even one not in subset (apache/datafusion#25829). The limitation was only in the PR description. Document it and the quoted-rename workaround in the missing values guide, add the missing fill_nan section there, and point to the note from both docstrings. Co-Authored-By: Claude Opus 5.5 (1M context) <[email protected]>
A positional filter now lands in mean's new distinct argument and still raises, as the paragraph above describes for the other functions. Co-Authored-By: Claude Opus 5.5 (1M context) <[email protected]>
f6d0478 fixed the _PyCapsule alias in user_defined.py, but context.py and extensions.py still imported CapsuleType from _typeshed, which does not define it, so the capsule half of set_query_planner and the codec methods was Unknown and accepted any argument. Import it from typing_extensions under TYPE_CHECKING there; type checkers ship its stubs on every Python version. In user_defined.py the alias existed only under TYPE_CHECKING, so typing.get_type_hints raised NameError on the from_pycapsule methods, two of which resolved on main. Import it at runtime; typing_extensions is already a runtime dependency below 3.13. Checked with pyright 1.1.414 at --pythonversion 3.10 and 3.13: the two unknown-import errors are gone, no new errors appear, and set_query_planner(1) and with_logical_extension_codec(1) are rejected. Co-Authored-By: Claude Opus 5.5 (1M context) <[email protected]>
Co-Authored-By: Claude Opus 5.5 (1M context) <[email protected]>
typing-extensions is a dependency only below 3.13, so context and extensions now pick the module the same way user_defined does. Co-Authored-By: Claude Opus 5.5 (1M context) <[email protected]>
6808f95 moved the CapsuleType import in user_defined.py out of TYPE_CHECKING so typing.get_type_hints can resolve the from_pycapsule overloads. typing_extensions only gained CapsuleType in 4.12.0, and the dependency had no floor, so an environment that pins an older release failed at `import datafusion` with ImportError. The [tool.uv] constraint only shapes this repo's lockfile and is not written into the wheel's metadata. Co-Authored-By: Claude Fable 5.1 <[email protected]>
_series read stop with bound.arguments.get, which returns None both
when the argument is omitted and when the caller passes None, so
range(col("n"), None) silently promoted n to the upper bound and
returned 0..n. BoundArguments only records arguments that were passed,
so check for presence and raise the same way the lone start= keyword
already does.
Co-Authored-By: Claude Fable 5.1 <[email protected]>
Passing order_by=[] to over() or to a window function keyword reached the builder as Some([]), so build() derived RANGE UNBOUNDED PRECEDING .. CURRENT ROW with no sort key. That expression fails at execution with "ORDER BY column cannot be empty". b4cad2d taught builder_from_expr to accept that frame as a default when chaining, but left the expression itself unexecutable. Filter the empty list out in apply_window_options so build() derives the whole-partition frame, and drop the special case and the doc sentence describing it. Add a test that order_by=[] executes. Co-Authored-By: Claude Fable 5.1 <[email protected]>
8027fcb inserted distinct before filter in percentile_cont, quantile_cont, mean, bit_and, bit_or, and string_agg. Only string_agg checked for a shifted positional call; the other five failed inside PyO3 with "'Expr' object is not an instance of 'bool'", which says nothing about which function or what to change. Move the string_agg check into a shared _check_distinct helper that names the function and the arguments to pass by keyword, and call it from all six. Co-Authored-By: Claude Fable 5.1 <[email protected]>
_series bound against a signature whose stop was optional and then needed three guards to undo that: start= alone, an explicit None stop, and step without stop. Each shape that fell through, such as range(stop=5) reporting a missing start, needed another. Model it as numpy.arange does: stop is required, start defaults to 0, and a lone positional argument is stop. range(stop=5) now works. range(start=5) and range(1, step=2) fail in Signature.bind with "missing a required argument: 'stop'", so the custom guards go. The one-argument overloads drop the positional-only marker so stop= type checks. Co-Authored-By: Claude Fable 5.1 <[email protected]>
The capsule branches of ScalarUDF, AggregateUDF, and WindowUDF from_pycapsule each rebuilt the "construct without __init__" idiom with object.__new__ and a cast, next to a _from_internal classmethod that already does exactly that. Call it instead so the three wrappers cannot drift from the registry lookups in SessionContext. Co-Authored-By: Claude Fable 5.1 <[email protected]>
PyExprFuncBuilder::from_expr matched the expression to find its kind and name, then called builder_from_expr, which matched the same expression again with the same three arms to seed the options. A variant handled in one and not the other would give a builder whose kind checks disagree with its options. Fold builder_from_expr into from_expr so each arm produces the builder, kind, and name together. over() takes the builder field from the result. Co-Authored-By: Claude Fable 5.1 <[email protected]>
The class docstring said each method raises when its option does not apply, but order_by and null_treatment never do; only filter, distinct, partition_by, and window_frame check the function kind. Keep the sentence on which options apply to which kind and leave the raising behavior to the upgrade guide, which already describes it. Co-Authored-By: Claude Fable 5.1 <[email protected]>
Remove tests that only pin upstream error messages for bad input (non-function builder, file metadata outside a scan, unknown fill_nan column, a bind error prefix), a type-hint accessor check, and cases that repeat an identical code path: the empty window order_by now normalizes to none so its rows duplicate the no-order_by rows, the approx percentile takes the same within-group branch, first_value takes the same aggregate arm as array_agg, the SQL window DISTINCT reaches the same accumulator check, and partition_by/window_frame and distinct/filter each share one kind check. Fold the remaining argument-shape checks for range into one parametrized test against range alone, since gen_series and generate_series share _series. Merge the string_agg and positional filter guards into one table, the two aggregate-window option tests into one, and the empty and numpy subset tests for fill_null and fill_nan into one each. Co-Authored-By: Claude Fable 5.1 <[email protected]>
A built window function stores a concrete frame with no record of whether it was chosen or derived from order_by, so merging a chained option into it had to guess. Each guess needed its own fix, and the last one dropped the sort key a decoded RANGE frame needs while keeping the frame, giving "ORDER BY column cannot be empty". Raise instead when a chain or over() starts from a window function with a partition_by, an order_by, or a frame other than the whole-partition default. null_treatment is kept, as are filter and distinct on an aggregate used as a window function. Aggregates still merge, since their options are stored as given. Merging can come back once apache/datafusion#25934 records whether a frame was explicit. Co-Authored-By: Claude Opus 5.5 <[email protected]>
The Null Treatment section said null treatment needs the builder, but its own example passes it in the Window, and over() now keeps one set on the aggregate. The higher-order functions list in expressions.md also missed array_first. Co-Authored-By: Claude Opus 5.5 <[email protected]>
…flag Co-Authored-By: Claude Opus 5.5 <[email protected]>
|
@andygrove Thanks again for the reviews. I think this is ready. My last run through with my agent didn't pick up anything new. |


Which issue does this PR close?
Closes #1764.
The other gaps were found by auditing the Python API against upstream DataFusion 55.1.0 with the
check-upstreamskill. Gaps that already have open issues (#1571–#1575, #1577, #1668, #1669) are left out.Rationale for this change
Several upstream functions, optional arguments, and DataFrame methods were not reachable from Python. The audit also turned up bugs where options on aggregate and window functions were silently dropped.
What changes are included in this PR?
array_first,any_value, file metadata functions, and a batch of Spark functions and pyspark aliases.distinct=on several aggregates,characters=on the trim functions,null_treatment=onlead/lag, and newDataFrame.explainoptions.explainraisesValueErrorfor option combinations the chosen plan would ignore (show_statisticswithanalyze, oranalyze_level/analyze_categorieswithout it), matching upstream SQL.DataFrame.fill_nan.ScalarUDFandWindowUDFaccept a bare PyCapsule, and the capsule overloads type-check correctly inuser_defined,context, andextensions, and resolve withtyping.get_type_hints. Passing the wrong kind of capsule toudf,udaf, orudwfnow names the expected and the actual capsule.rangeandgen_series: they accept plain ints, andrange(stop)works like Python's built-in. The lone upper bound is taken only by position, sorange(start=5)raises instead of meaningstop=5..over()onto an aggregate or window function no longer discards options already set on it (ordering, null treatment, window frame). Whether a window frame is re-derived depends only on the expression, so it behaves the same aftercopy,pickle,from_bytes, or parsing from SQL.filteranddistinctneed an aggregate (including one used as a window function), andpartition_byandwindow_frameneed a window function..over()on an aggregate keeps itsfilter,distinct, andnull_treatment(Expr.over() on an aggregate silently drops order_by, null_treatment, filter, and distinct #1764). An aggregateorder_byraises, asORDER BYinside an aggregate call does withOVERin SQL; aWITHIN GROUPfunction such aspercentile_contaccepts an ascendingsort_expressionand raises on a descending one.percentile_cont,quantile_cont,approx_percentile_cont, andapprox_percentile_cont_with_weightkeep the direction of their sort expression.DISTINCTinstead of silently counting every row. A query the optimizer rewrites to group by the distinct values first still runs.mean(x, filter=...)no longer raisesTypeError.string_aggrejectsNoneor anExprasdistinct, so an old positionalorder_bycall fails loudly instead of being used as the filter.numpy.bool_is still accepted.fill_null(subset=[])andfill_nan(subset=[])fill no columns instead of all of them. numpy and pandas column lists still work.spark.printftreats a barestras a column name, as pyspark does.udf,udaf, andudwfdecorators no longer raisesIndexError, and the function form accepts the callable by keyword (udf(func=...),udaf(accum=...),udwf(func=...),udtf(func=...)).substrguidance inskills/datafusion_python/SKILL.mdnow shows thelengthargument as valid.:ref:roles.count_stardocstrings describe what the builder options actually do.fill_nanandfill_nullfail on uppercase or dotted column names because of an upstream bug, filed as DataFrame::fill_null and fill_nan fail on column names with uppercase letters or dots datafusion#25829. The functions guide documents the limitation and a rename workaround, and both docstrings point to it.Every new function has a doctest and pytest coverage.
Are there any user-facing changes?
Yes, new functions, arguments, and methods. The following are documented in
docs/source/user-guide/upgrade-guides.md:distinctis inserted beforefilterinbit_and,bit_or,mean,percentile_cont,quantile_cont, andstring_agg, matchingsumandavgin 54.0.0. Passfilter(and, forstring_agg,order_by) by keyword.array_agg(..., distinct=True).order_by(...)), and generated names forfirst_value,last_value, andnth_valuenow includeRESPECT NULLS..over()keeps an aggregate'sfilter,distinct, andnull_treatment, which can change results. An aggregateorder_bynow raises; move it into theWindow.DISTINCTinstead of returning the non-distinct result.fill_null(subset=[])fills no columns.spark.last_day's parameter is renamed fromcoltodateto match pyspark.Not in the upgrade guide, because each only turns a silently wrong or ignored call into an error:
explainrejects option combinations the plan would ignore, andrange(start=...)/gen_series(start=...)withoutstopraise.🤖 Generated with Claude Code