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 <noreply@anthropic.com>
…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 <noreply@anthropic.com>
- 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 <noreply@anthropic.com>
…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 <noreply@anthropic.com>
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 <noreply@anthropic.com>
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 <noreply@anthropic.com>
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 <noreply@anthropic.com>
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 <noreply@anthropic.com>
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 <noreply@anthropic.com>
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 <noreply@anthropic.com>
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 <noreply@anthropic.com>
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 <noreply@anthropic.com>
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 <noreply@anthropic.com>
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 <noreply@anthropic.com>
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 <noreply@anthropic.com>
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 <noreply@anthropic.com>
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.
| """Graphviz DOT format for graph rendering.""" | ||
|
|
||
|
|
||
| class ExplainAnalyzeLevel(Enum): |
There was a problem hiding this comment.
Not taking this one. The package root has grown to hold nearly every public class, which makes it hard to predict where anything lives: ExplainFormat is at the root, and so is WindowFrame, but Window is not. 55.0.0 already started the other direction: the new extension protocols live only in datafusion.extensions, and SessionExtensionComponents stays at the root "because a bundle constructs one rather than merely naming it". ExplainAnalyzeLevel and ExplainMetricCategory are only ever named as arguments to DataFrame.explain, so they stay in datafusion.dataframe. I've opened #1771 to apply the same rule to the rest of the root, with a deprecation period, rather than extend the current pattern here.
| Returns NULL if every value in the group is NULL. Which value is returned | ||
| is not specified and may differ between runs. | ||
|
|
||
| If using the builder functions described in ref:`_aggregation` this function ignores |
There was a problem hiding this comment.
Fixed in 7259819: docs: fix the broken guide links in aggregate docstrings
| ("bit_and_distinct", f.bit_and(column("b"), distinct=True), [4]), | ||
| ("bit_or_distinct", f.bit_or(column("b"), distinct=True), [6]), |
There was a problem hiding this comment.
Fixed in d779502: test: check that bit_and and bit_or keep distinct in the expression
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.
| /// chose it is lost. `keep_window_frame` carries that from the Python side; when | ||
| /// false, a frame equal to the default for the current order-by is treated as | ||
| /// unset. | ||
| fn builder_from_expr(expr: &Expr, keep_window_frame: bool) -> ExprFuncBuilder { |
There was a problem hiding this comment.
Because the seeded builder always carries a function kind, upstream's per-method kind checks in ExprFunctionExt no longer apply. partition_by/window_frame on a plain aggregate, and distinct on a non-aggregate window function, now build and the option is silently dropped. On main each of these raises at build() with "ExprFunctionExt can only be used with Expr::AggregateFunction or Expr::WindowFunction":
from datafusion import SessionContext, col, functions as f
from datafusion.expr import WindowFrame
ctx = SessionContext()
df = ctx.from_pydict({"g": [1, 1, 1, 2], "v": [3, 1, 2, 4]})
e = f.sum(col("v")).partition_by(col("g")).build()
print(e) # Expr(sum(v))
print(df.aggregate([], [e.alias("s")]).to_pydict()) # {'s': [10]}
print(f.sum(col("v")).window_frame(WindowFrame("rows", 1, 0)).build()) # Expr(sum(v))
print(f.lead(col("v"), order_by="v").distinct().build())
# Expr(lead(DISTINCT v, Int64(1), NULL) ORDER BY [...]) -- runs, DISTINCT ignoredThe new filter/distinct support on aggregate windows (f.sum(..).over(..).distinct()) does work, so one way to keep it while restoring the errors: in partition_by/window_frame, fall back to upstream's self.expr.clone().partition_by(..) / .window_frame(..) when self.expr is an AggregateFunction, and in distinct do the same for a window function whose definition isn't WindowFunctionDefinition::AggregateUDF.
There was a problem hiding this comment.
Fixed in 92ace16: fix: raise when a builder option does not apply to the function kind
| def string_agg( | ||
| expression: Expr, | ||
| delimiter: str, | ||
| distinct: bool = False, |
There was a problem hiding this comment.
The upgrade guide covers the distinct insertion, but this one positional shape doesn't fail loudly like the others do. string_agg(expr, ",", None, order_col) now binds None to distinct (the Rust Option<bool> accepts it) and the ordering column to filter:
from datafusion import SessionContext, col, functions as f
ctx = SessionContext()
df = ctx.from_pydict({"s": ["c", "a", "b", "d"], "n": [3, 1, 0, 2]})
e = f.string_agg(col("s"), ",", None, col("n"))
print(df.aggregate([], [e.alias("r")]).to_pydict())
# main: {'r': ['b,a,d,c']} (ORDER BY n)
# PR: {'r': ['c,a,d']} (n used as the FILTER)The other shifted calls raise TypeError: 'Expr' object is not an instance of 'bool'. A guard such as if not isinstance(distinct, bool): raise TypeError(...) would make this one fail too.
There was a problem hiding this comment.
Fixed in 0ac494e: fix: reject a non-bool distinct in string_agg
| functions_aggregate::expr_fn::percentile_cont(sort_expression.sort, lit(percentile)); | ||
|
|
||
| add_builder_fns_to_aggregate(agg_fn, None, filter, None, None) | ||
| add_builder_fns_to_aggregate(agg_fn, distinct, filter, None, None) |
There was a problem hiding this comment.
Pre-existing, but since this call is being edited: add_builder_fns_to_aggregate starts from an empty builder (agg_fn.null_treatment(None)), so build() overwrites the WITHIN GROUP ordering that upstream's percentile_cont(sort, p) stored. A descending sort is silently ignored:
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]}, name="t")
e = f.percentile_cont(col("a").sort(ascending=False), 0.25)
print(e) # Expr(percentile_cont(a, Float64(0.25))), no DESC
print(df.aggregate([], [e.alias("p")]).to_pydict()) # {'p': [2.0]}
print(ctx.sql(
"SELECT percentile_cont(0.25) WITHIN GROUP (ORDER BY a DESC) AS p FROM t"
).to_pydict()) # {'p': [4.0]}quantile_cont, approx_percentile_cont (1.75 vs 4.25 from SQL) and approx_percentile_cont_with_weight go through the same path. Passing the sort expression through as order_by here would keep it. Relatedly, the docstring's "this function ignores the options order_by and null_treatment" isn't right for order_by: f.percentile_cont(col("a"), 0.25).order_by(col("a").sort(ascending=False)).build() returns 4.0.
There was a problem hiding this comment.
Fixed in 1b01845: fix: keep the sort direction in percentile_cont and friends
| }; | ||
|
|
||
| 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.
| ctx.execute(plan, partition=0) # after | ||
| ``` | ||
|
|
||
| ### More aggregate functions accept `distinct` |
There was a problem hiding this comment.
The PR description says chaining builder methods or .over() onto an already-configured function may now give different results, but that isn't in the upgrade guide yet. Some examples of what changes:
from datafusion import SessionContext, col, functions as f
from datafusion.expr import Window
ctx = SessionContext()
df = ctx.from_pydict(
{"g": [1, 1, 1, 2, 2], "t": [1, 2, 3, 4, 5], "v": [10, None, 30, 40, 50]}
)
e = f.lead(col("v"), 1, partition_by=[col("g")], order_by="t").over(Window(order_by="t"))
print(df.select(col("t"), e.alias("r")).sort(col("t")).collect_column("r").to_pylist())
# main: [None, 30, 40, 50, None]
# PR: [None, 30, None, 50, None] (partition now kept)
df2 = ctx.from_pydict({"s": ["a", "b", "a"], "v": [3, 1, 2]})
e2 = f.array_agg(col("s"), distinct=True).order_by(col("v")).build()
print(df2.aggregate([], [e2.alias("r")]).to_pydict())
# main: {'r': [['b', 'a', 'a']]} (distinct silently dropped)
# PR: Execution error: In an aggregate with DISTINCT, ORDER BY expressions must appear in argument listAuto-generated names change too. f.first_value(col("a")).order_by(col("b")).build(), the chained form shown in aggregations.md, was named first_value(a) ORDER BY [b ASC NULLS FIRST] on main and is now first_value(a) RESPECT NULLS ORDER BY [b ASC NULLS FIRST]. That matches the non-chained form, but code that selects the unaliased column by name will break. last_value and nth_value change the same way. Expr.over's docstring also still describes only aggregates.
There was a problem hiding this comment.
Fixed in 6b52225: docs: document that chaining now keeps options already set
| analyze: If ``True``, the plan will run and metrics reported. | ||
| format: Output format for the plan. Defaults to | ||
| :py:attr:`ExplainFormat.INDENT`. | ||
| show_statistics: If ``True``, include each operator's statistics. |
There was a problem hiding this comment.
show_statistics is ignored when analyze=True. Upstream's LogicalPlan::Analyze has statement-level overrides for analyze_level and analyze_categories, but none for statistics:
from datafusion import SessionContext, col, lit
df = SessionContext().from_pydict({"a": [1, 2]}).filter(col("a") > lit(1))
df.explain(show_statistics=True) # physical plan includes statistics=[...]
df.explain(analyze=True, show_statistics=True) # no statisticsUpstream SQL rejects the combination (EXPLAIN option COSTS cannot be combined with ANALYZE), so raising a ValueError here, or at least saying so in the docstring, would match.
There was a problem hiding this comment.
Fixed in 607dde4: fix: reject show_statistics combined with analyze in explain
| return Expr(_f.format_string(fmt_expr.expr, *[c.expr for c in cols])) | ||
|
|
||
|
|
||
| def printf(format: str | Expr, *cols: Expr) -> Expr: |
There was a problem hiding this comment.
pyspark's printf(format: ColumnOrName, *cols) treats a bare str as a column name, unlike format_string(format: str, ...). As an alias of format_string, this silently returns the column name's text for pyspark-style calls:
from datafusion import SessionContext, col
from datafusion.functions import spark
df = SessionContext().from_pydict({"a": ["aa%d%s"], "b": [123], "c": ["cc"]})
print(df.select(spark.printf("a", col("b"), col("c")).alias("v")).to_pydict())
# {'v': ['a']}; pyspark's documented sf.printf("a", "b", "c") gives 'aa123cc'
print(df.select(spark.printf(col("a"), col("b"), col("c")).alias("v")).to_pydict())
# {'v': ['aa123cc']}Other functions in this module already treat a bare str as a column name to match pyspark (e.g. bit_get's pos).
There was a problem hiding this comment.
Fixed in 4933fa2: fix: treat a bare str as a column name in spark.printf
| } | ||
| } | ||
|
|
||
| fn scalar_udf_from_capsule(capsule: &Bound<'_, PyCapsule>) -> PyDataFusionResult<ScalarUDF> { |
There was a problem hiding this comment.
Now that a bare capsule is accepted, passing the wrong kind is easy, and pointer_checked on its own only gives CPython's fixed message:
from datafusion import SessionContext, udf
ctx = SessionContext()
udf(ctx.__datafusion_logical_extension_codec__())
# ValueError: PyCapsule_GetPointer called with incorrect name
ctx.register_table("t", ctx.__datafusion_logical_extension_codec__())
# ValueError: Expected name 'datafusion_table_provider' in PyCapsule, instead got 'datafusion_logical_extension_codec'validate_pycapsule in crates/util/src/lib.rs exists for this, and its doc comment says to call it first at every extraction site. Adding validate_pycapsule(capsule, "datafusion_scalar_udf")?; here, plus the equivalent in window_udf_from_capsule and aggregate_udf_from_capsule, would fix it.
There was a problem hiding this comment.
Fixed in 23ed533: fix: name the expected and found capsule when importing a UDF
|
|
||
| Args: | ||
| value: Value to replace NaN with. Will be cast to match column type. | ||
| subset: Optional list of column names to fill. If None, fills all |
There was a problem hiding this comment.
subset=[] fills every floating-point column, not none. The binding maps both None and [] to an empty Vec, which upstream treats as "all columns":
from datafusion import SessionContext
nan = float("nan")
df = SessionContext().from_pydict({"a": [1.0, nan, None], "b": [nan, 2.0, 3.0]})
print(df.fill_nan(0.0, subset=[]).to_pydict())
# {'a': [1.0, 0.0, None], 'b': [0.0, 2.0, 3.0]}A computed subset that happens to be empty ([c for c in cols if pred(c)]) would rewrite every float column. Returning self unchanged for [], or documenting it, would avoid the surprise. fill_null has the same behavior.
There was a problem hiding this comment.
Fixed in 8454a09: fix: treat an empty fill_null or fill_nan subset as no columns
| return decorator | ||
|
|
||
| if hasattr(args[0], "__datafusion_scalar_udf__"): | ||
| if hasattr(args[0], "__datafusion_scalar_udf__") or _is_pycapsule(args[0]): |
There was a problem hiding this comment.
Pre-existing (main behaves the same, and udaf too), but since this line changed: args[0] is read before the if args and ... guard below, so the keyword-only decorator form raises IndexError:
import pyarrow as pa
from datafusion import udf, udwf
udf(input_fields=[pa.int64()], return_field=pa.int64(), volatility="immutable")
# IndexError: tuple index out of range
udwf(input_types=[pa.int64()], return_type=pa.int64(), volatility="immutable")
# IndexError: tuple index out of rangeif args and (hasattr(args[0], "__datafusion_scalar_udf__") or _is_pycapsule(args[0])): fixes it, and the same at line 1101 for udwf.
There was a problem hiding this comment.
Fixed in f8d7b70: fix: allow the keyword-only form of the udf, udaf, and udwf decorators
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) <noreply@anthropic.com>
…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) <noreply@anthropic.com>
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) <noreply@anthropic.com>
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) <noreply@anthropic.com>
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) <noreply@anthropic.com>
Inserting distinct before filter shifts positional arguments. Most shifted calls fail loudly because an Expr is not a bool, but the old way to pass order_by positionally, string_agg(expr, ",", None, order), still ran: the Rust binding takes Option<bool>, so None became distinct=False and the ordering column became the filter. Raise a TypeError for a non-bool distinct that says to pass filter and order_by by keyword. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
Upstream's LogicalPlan::Analyze has statement-level overrides for the
analyze level and categories but none for statistics, so
explain(analyze=True, show_statistics=True) silently printed no
statistics. SQL rejects the same combination ("EXPLAIN option COSTS
cannot be combined with ANALYZE"); raise a ValueError to match.
Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
pyspark's printf(format: ColumnOrName, *cols: ColumnOrName) reads a bare
str as a column name, unlike format_string, whose format is a plain str
template. As an alias of format_string, spark.printf("a", ...) returned
the literal text of the column name instead of formatting with it.
Give printf its own body that resolves str arguments as column names, as
bit_get already does for pos.
Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
Now that udf, udaf, and udwf accept a bare PyCapsule, passing the wrong kind is easy, and pointer_checked alone surfaces CPython's fixed "PyCapsule_GetPointer called with incorrect name". Call validate_pycapsule first in the scalar, window, and aggregate capsule imports, as its doc comment asks of every extraction site, so the error names both the expected and the actual capsule. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
Both bindings collect subset into a Vec, and upstream reads an empty column list as "all columns", so subset=[] behaved like subset=None. A subset computed from the schema that happened to match nothing then rewrote every column (every float column for fill_nan). Return the DataFrame unchanged for an empty subset, matching pyspark's na.fill, and note the fill_null change in the upgrade guide. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
Upstream's percentile_cont(sort, p) stores the sort as the aggregate's WITHIN GROUP ordering, but add_builder_fns_to_aggregate starts from an empty builder, so build() overwrote it. A descending sort was silently ignored in percentile_cont, quantile_cont, approx_percentile_cont, and approx_percentile_cont_with_weight (2.0 instead of SQL's 4.0 for the 0.25 percentile of 1..5). Pass the sort through as order_by. Generated column names now include the ordering, as they do from SQL; note that in the upgrade guide. Correct the docstrings, which said these functions ignore order_by: a chained order_by replaces the ordering of sort_expression. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
The capsule check read args[0] before the "if args" guard that follows it, so calling a decorator with only keyword arguments, such as udf(input_fields=..., return_field=..., volatility=...), raised IndexError instead of returning the decorator. Guard the capsule check on args being non-empty, as udtf already does. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
The aggregate docstrings linked the aggregation guide as ref:`_aggregation`, with no leading colon and a leading underscore on the label, so it rendered as plain italic text instead of a link. The pattern came in with the builder parameters and was copied into each new aggregate since; the NullTreatment docstring and the window function note had the same mistake with _window_functions. Use :ref:`aggregation` and :ref:`window_functions`, matching the labels in the guide and the module docstring that already links them. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
AND and OR ignore duplicate inputs, so the existing result checks for bit_and(distinct=True) and bit_or(distinct=True) pass whether or not the argument reaches the plan. Assert on the built expression's canonical name as well. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
The skill said F.substr takes only two arguments and labelled
F.substr(col("c_phone"), lit(1), lit(2)) as wrong. substr now accepts an
optional length, so that call returns the first two characters, and the
skill was steering agents away from a valid form.
Describe both forms, mark length as requiring datafusion-python 55, and
keep the TypeError note for older versions.
Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com>
|
@andygrove re the |
|
@andygrove Thanks for the review! I think we've addressed all of the comments, Claude and I ;) |


Which issue does this PR close?
No issue filed. The 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.explainraisesValueErrorwhenshow_statisticsis combined withanalyze, matching upstream SQL.DataFrame.fill_nan.ScalarUDFandWindowUDFaccept a bare PyCapsule, and the capsule overloads now type-check correctly. 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..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, or parsing from SQL. A builder option that does not apply to the function's kind still raises.percentile_cont,quantile_cont,approx_percentile_cont, andapprox_percentile_cont_with_weightkeep the direction of their sort expression.mean(x, filter=...)no longer raisesTypeError.string_aggrejects a non-booldistinct, so an old positionalorder_bycall fails loudly instead of being used as the filter.fill_null(subset=[])andfill_nan(subset=[])fill no columns instead of all of them.spark.printftreats a barestras a column name, as pyspark does.udf,udaf, andudwfdecorators no longer raisesIndexError.substrguidance inskills/datafusion_python/SKILL.mdnow shows thelengthargument as valid, and the aggregate docstrings link to the guide with real:ref:roles.Every new function has a doctest and pytest coverage. The same options are still dropped when
.over()converts an aggregate to a window function; that is tracked separately in #1764 and noted in theover()docstring.fill_nanandfill_nullfail on uppercase or dotted column names because of an upstream bug, filed as apache/datafusion#25829.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.fill_null(subset=[])fills no columns.spark.last_day's parameter is renamed fromcoltodateto match pyspark.🤖 Generated with Claude Code