Skip to content

feat: close upstream coverage gaps for DataFusion 55.1.0 - #1763

Open
timsaucer wants to merge 31 commits into
apache:mainfrom
timsaucer:feat/upstream-coverage-gaps
Open

timsaucer wants to merge 31 commits into
apache:mainfrom
timsaucer:feat/upstream-coverage-gaps

Conversation

@timsaucer

@timsaucer timsaucer commented Sep 25, 2026 •

Copy link
Copy Markdown
Member

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-upstream skill. 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?

  • New functions: array math and array_first, any_value, file metadata functions, and a batch of Spark functions and pyspark aliases.
  • Optional arguments upstream already supports, such as distinct= on several aggregates, characters= on the trim functions, null_treatment= on lead/lag, and new DataFrame.explain options. explain raises ValueError when show_statistics is combined with analyze, matching upstream SQL.
  • DataFrame.fill_nan.
  • FFI: ScalarUDF and WindowUDF accept a bare PyCapsule, and the capsule overloads now type-check correctly. Passing the wrong kind of capsule to udf, udaf, or udwf now names the expected and the actual capsule.
  • More Pythonic range and gen_series: they accept plain ints, and range(stop) works like Python's built-in.
  • Fixes:
    • Chaining builder methods or .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 after copy, 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, and approx_percentile_cont_with_weight keep the direction of their sort expression.
    • mean(x, filter=...) no longer raises TypeError.
    • string_agg rejects a non-bool distinct, so an old positional order_by call fails loudly instead of being used as the filter.
    • fill_null(subset=[]) and fill_nan(subset=[]) fill no columns instead of all of them.
    • spark.printf treats a bare str as a column name, as pyspark does.
    • The keyword-only form of the udf, udaf, and udwf decorators no longer raises IndexError.
  • Docs: the substr guidance in skills/datafusion_python/SKILL.md now shows the length argument 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 the over() docstring. fill_nan and fill_null fail 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:

  • distinct is inserted before filter in bit_and, bit_or, mean, percentile_cont, quantile_cont, and string_agg, matching sum and avg in 54.0.0. Pass filter (and, for string_agg, order_by) by keyword.
  • Chaining keeps options already set. Results can change, some chains that ran before now raise (for example array_agg(..., distinct=True).order_by(...)), and generated names for first_value, last_value, and nth_value now include RESPECT NULLS.
  • The percentile functions honor a descending sort, and their generated names include the ordering.
  • fill_null(subset=[]) fills no columns.
  • spark.last_day's parameter is renamed from col to date to match pyspark.

🤖 Generated with Claude Code

timsaucer and others added 6 commits September 25, 2026 08:19
…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>
@timsaucer
timsaucer force-pushed the feat/upstream-coverage-gaps branch from 5f33af0 to 7cf8fb5 Compare September 25, 2026 13:22
timsaucer and others added 5 commits September 25, 2026 09:44
…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>
@timsaucer
timsaucer force-pushed the feat/upstream-coverage-gaps branch from 7cf8fb5 to 5d43576 Compare September 25, 2026 13:45
timsaucer and others added 5 commits September 25, 2026 10:53
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>
@timsaucer
timsaucer marked this pull request as ready for review September 25, 2026 16:30
@timsaucer
timsaucer requested a balanced review from Copilot September 25, 2026 16:30
@timsaucer timsaucer changed the title feat: close upstream coverage gaps found by audit against DataFusion 55.1.0 feat: close upstream coverage gaps for DataFusion 55.1.0 Sep 25, 2026

Copilot AI left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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 Medium severity · 2 Low severity

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):

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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.

Comment thread python/datafusion/functions/__init__.py Outdated
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

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Fixed in 7259819: docs: fix the broken guide links in aggregate docstrings

Comment on lines +367 to +368
("bit_and_distinct", f.bit_and(column("b"), distinct=True), [4]),
("bit_or_distinct", f.bit_or(column("b"), distinct=True), [6]),

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Fixed in d779502: test: check that bit_and and bit_or keep distinct in the expression

@andygrove andygrove left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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.

Comment thread crates/core/src/expr.rs Outdated
/// 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 {

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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 ignored

The 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.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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,

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Fixed in 0ac494e: fix: reject a non-bool distinct in string_agg

Comment thread crates/core/src/functions.rs Outdated
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)

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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)?;

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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`

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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 list

Auto-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.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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.

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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 statistics

Upstream 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.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Fixed in 607dde4: fix: reject show_statistics combined with analyze in explain

Comment thread python/datafusion/functions/spark.py Outdated
return Expr(_f.format_string(fmt_expr.expr, *[c.expr for c in cols]))


def printf(format: str | Expr, *cols: Expr) -> Expr:

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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).

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Fixed in 4933fa2: fix: treat a bare str as a column name in spark.printf

Comment thread crates/core/src/udf.rs
}
}

fn scalar_udf_from_capsule(capsule: &Bound<'_, PyCapsule>) -> PyDataFusionResult<ScalarUDF> {

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Fixed in 8454a09: fix: treat an empty fill_null or fill_nan subset as no columns

Comment thread python/datafusion/user_defined.py Outdated
return decorator

if hasattr(args[0], "__datafusion_scalar_udf__"):
if hasattr(args[0], "__datafusion_scalar_udf__") or _is_pycapsule(args[0]):

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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 range

if args and (hasattr(args[0], "__datafusion_scalar_udf__") or _is_pycapsule(args[0])): fixes it, and the same at line 1101 for udwf.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Fixed in f8d7b70: fix: allow the keyword-only form of the udf, udaf, and udwf decorators

timsaucer and others added 5 commits September 28, 2026 08:02
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>
timsaucer and others added 7 commits September 28, 2026 08:20
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>
timsaucer and others added 3 commits September 28, 2026 12:04
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>
@timsaucer

Copy link
Copy Markdown
Member Author

@andygrove re the SKILL.md substr note in the review summary: fixed in db5b7cc: docs(skill): show substr's length argument as valid

@timsaucer

Copy link
Copy Markdown
Member Author

@andygrove Thanks for the review! I think we've addressed all of the comments, Claude and I ;)

@andygrove andygrove left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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_cont and friends keep the direction as aggregates, but .over() still drops it, and the windows.md workaround can't bring it back (1b01845),
  • the empty-order_by frame (b4cad2d) is kept after copy/pickle/from_bytes, because decoding adds ORDER BY UInt64(1),
  • the empty-subset guard (8454a09) raises on numpy arrays and pandas.Index subsets, which worked on main,
  • the kind check (92ace16) only covers the first builder call, and filter() has none,
  • explain (607dde4) still ignores show_statistics=False with analyze, and analyze_level/analyze_categories without it,
  • the #1764 workaround in windows.md (8e21c8a) gives the non-distinct result for Python UDAFs,
  • string_agg's bool check (0ac494e) rejects np.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.

Comment thread crates/core/src/expr.rs
@@ -678,8 +705,8 @@ impl PyExpr {
null_treatment,

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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.

Comment thread crates/core/src/expr.rs
// 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 {

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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.

- 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:

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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.

Comment thread crates/core/src/expr.rs
.into()
}

pub fn filter(&self, filter: PyExpr) -> PyExprFuncBuilder {

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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.

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

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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.

>>> result.collect_column("s")[0].as_py()
'x,y'
"""
if not isinstance(distinct, bool):

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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 keyword

A 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.

):
return ScalarUDF.from_pycapsule(args[0])

if args and callable(args[0]):

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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.

"""
if subset is not None and not subset:
return self
return DataFrame(self.df.fill_nan(value, subset))

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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"'?

f.bit_and(column("a"), filter=my_filter) # after
```

Passing `filter` to `mean` previously raised a `TypeError`, whether passed

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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.

import sys

if sys.version_info >= (3, 13):
from types import CapsuleType as _PyCapsule

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

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.

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Projects

None yet

Development

Successfully merging this pull request may close these issues.

3 participants