andygrove commented on code in PR #1763:
URL:
https://github.com/apache/datafusion-python/pull/1763#discussion_r4115775488
##########
crates/core/src/expr.rs:
##########
@@ -743,6 +767,63 @@ impl PyExpr {
}
}
+/// Start an [`ExprFuncBuilder`] that keeps the options already set on `expr`.
+///
+/// Upstream's `ExprFunctionExt` methods on an `Expr` start from an empty
+/// builder, so `build()` would reset every option not set again. The Python
+/// function wrappers already apply their keyword options, so chaining another
+/// builder method onto their result must not discard them.
+///
+/// A built window function always stores a concrete frame, so whether the user
+/// 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 {
Review 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":
```python
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`.
##########
python/datafusion/functions/__init__.py:
##########
@@ -7054,6 +7576,7 @@ def ntile(
def string_agg(
expression: Expr,
delimiter: str,
+ distinct: bool = False,
Review 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`:
```python
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.
##########
crates/core/src/dataframe.rs:
##########
@@ -1320,6 +1345,26 @@ impl PyDataFrame {
let df = self.df.as_ref().fill_null(&scalar_value.0, &cols)?;
Ok(Self::new(df))
}
+
+ /// Fill NaN values with a specified value for specific floating-point
columns
+ #[pyo3(signature = (value, columns=None))]
+ fn fill_nan(
+ &self,
+ value: Py<PyAny>,
+ columns: Option<Vec<PyBackedStr>>,
+ py: Python,
+ ) -> PyDataFusionResult<Self> {
+ let scalar_value: PyScalarValue = value.extract(py)?;
+
+ let cols = match columns {
+ Some(col_names) => col_names.iter().map(|c|
c.to_string()).collect(),
+ None => Vec::new(), // Empty vector means fill NaN for all columns
+ };
+
+ let cols = cols.iter().map(String::as_str).collect::<Vec<_>>();
+ let df = self.df.as_ref().fill_nan(&scalar_value.0, &cols)?;
Review 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`:
```python
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`.
##########
crates/core/src/expr.rs:
##########
@@ -624,44 +624,68 @@ impl PyExpr {
// Expression Function Builder functions
- pub fn order_by(&self, order_by: Vec<PySortExpr>) -> PyExprFuncBuilder {
- self.expr
- .clone()
+ #[pyo3(signature = (order_by, keep_window_frame=false))]
+ pub fn order_by(
+ &self,
+ order_by: Vec<PySortExpr>,
+ keep_window_frame: bool,
+ ) -> PyExprFuncBuilder {
+ builder_from_expr(&self.expr, keep_window_frame)
.order_by(to_sort_expressions(order_by))
.into()
}
- pub fn filter(&self, filter: PyExpr) -> PyExprFuncBuilder {
- self.expr.clone().filter(filter.expr.clone()).into()
+ #[pyo3(signature = (filter, keep_window_frame=false))]
+ pub fn filter(&self, filter: PyExpr, keep_window_frame: bool) ->
PyExprFuncBuilder {
+ builder_from_expr(&self.expr, keep_window_frame)
+ .filter(filter.expr.clone())
+ .into()
}
- pub fn distinct(&self) -> PyExprFuncBuilder {
- self.expr.clone().distinct().into()
+ #[pyo3(signature = (keep_window_frame=false))]
+ pub fn distinct(&self, keep_window_frame: bool) -> PyExprFuncBuilder {
+ builder_from_expr(&self.expr, keep_window_frame)
+ .distinct()
+ .into()
}
- pub fn null_treatment(&self, null_treatment: NullTreatment) ->
PyExprFuncBuilder {
- self.expr
- .clone()
+ #[pyo3(signature = (null_treatment, keep_window_frame=false))]
+ pub fn null_treatment(
+ &self,
+ null_treatment: NullTreatment,
+ keep_window_frame: bool,
+ ) -> PyExprFuncBuilder {
+ builder_from_expr(&self.expr, keep_window_frame)
.null_treatment(Some(null_treatment.into()))
.into()
}
- pub fn partition_by(&self, partition_by: Vec<PyExpr>) -> PyExprFuncBuilder
{
+ #[pyo3(signature = (partition_by, keep_window_frame=false))]
+ pub fn partition_by(
+ &self,
+ partition_by: Vec<PyExpr>,
+ keep_window_frame: bool,
+ ) -> PyExprFuncBuilder {
let partition_by = partition_by.iter().map(|e|
e.expr.clone()).collect();
- self.expr.clone().partition_by(partition_by).into()
+ builder_from_expr(&self.expr, keep_window_frame)
+ .partition_by(partition_by)
+ .into()
}
pub fn window_frame(&self, window_frame: PyWindowFrame) ->
PyExprFuncBuilder {
- self.expr.clone().window_frame(window_frame.into()).into()
+ builder_from_expr(&self.expr, false)
+ .window_frame(window_frame.into())
+ .into()
}
- #[pyo3(signature = (partition_by=None, window_frame=None, order_by=None,
null_treatment=None))]
+ #[pyo3(signature = (partition_by=None, window_frame=None, order_by=None,
null_treatment=None, keep_window_frame=false))]
pub fn over(
&self,
partition_by: Option<Vec<PyExpr>>,
window_frame: Option<PyWindowFrame>,
order_by: Option<Vec<PySortExpr>>,
null_treatment: Option<NullTreatment>,
+ keep_window_frame: bool,
) -> PyDataFusionResult<PyExpr> {
match &self.expr {
Expr::AggregateFunction(agg_fn) => {
Review Comment:
This is #1764, but the new `distinct=` arguments on `mean`,
`percentile_cont`, `quantile_cont` and `string_agg` run straight into it:
```python
from datafusion import SessionContext, col, functions as f
from datafusion.expr import Window
ctx = SessionContext()
df = ctx.from_pydict({"v": [1.0, 1.0, None, 4.0]}, name="t")
print(df.select(f.mean(col("v"),
distinct=True).over(Window()).alias("m")).to_pydict())
# {'m': [2.0, 2.0, 2.0, 2.0]}
print(ctx.sql("SELECT avg(DISTINCT v) OVER () AS m FROM t").to_pydict())
# {'m': [2.5, 2.5, 2.5, 2.5]}
```
Until #1764 lands, a note in the `over()` docstring (or on the new
`distinct` parameters) would help.
##########
python/datafusion/expr.py:
##########
@@ -453,6 +453,10 @@ class Expr: # noqa: PLW1641
:ref:`Expressions` in the online documentation for more information.
"""
+ # Set by ``over()`` when the window frame was given explicitly, so
chaining a
+ # builder method keeps it even if it equals the default frame.
+ _explicit_window_frame = False
Review Comment:
Because "the frame was explicit" lives only on the Python object, anything
that reconstructs the `Expr` loses it, and the same builder chain then gives
different results for identical expressions:
```python
import copy
import pickle
from datafusion import SessionContext, col, functions as f
from datafusion.expr import Window, WindowFrame
ctx = SessionContext()
df = ctx.from_pydict({"v": [3, 1, 2, 4]})
e = f.sum(col("v")).over(Window(window_frame=WindowFrame("rows", None,
None)))
def run(x):
r = x.order_by(col("v")).build().alias("r")
return df.select(col("v"),
r).sort(col("v")).collect_column("r").to_pylist()
print(run(e)) # [10, 10, 10, 10]
print(run(copy.copy(e))) # [1, 3, 6, 10]
print(run(pickle.loads(pickle.dumps(e)))) # [1, 3, 6, 10]
print(run(df.parse_sql_expr(
"sum(v) OVER (ROWS BETWEEN UNBOUNDED PRECEDING AND UNBOUNDED FOLLOWING)"
))) # [1, 3, 6, 10], though the frame
is written out
```
Pickling is how expressions get shipped to workers, so this can differ
between a driver and its workers. A rule that depends only on the expression
itself would behave the same everywhere, even if it's less precise.
##########
crates/core/src/udf.rs:
##########
@@ -209,6 +209,16 @@ impl ScalarUDFImpl for PythonFunctionScalarUDF {
}
}
+fn scalar_udf_from_capsule(capsule: &Bound<'_, PyCapsule>) ->
PyDataFusionResult<ScalarUDF> {
Review 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:
```python
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.
##########
python/datafusion/functions/spark.py:
##########
@@ -1520,6 +1711,15 @@ def format_string(format: str | Expr, *cols: Expr) ->
Expr:
return Expr(_f.format_string(fmt_expr.expr, *[c.expr for c in cols]))
+def printf(format: str | Expr, *cols: Expr) -> Expr:
Review 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:
```python
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`).
##########
python/datafusion/user_defined.py:
##########
@@ -388,7 +397,7 @@ def wrapper(*args: Any, **kwargs: Any) -> Callable:
return decorator
- if hasattr(args[0], "__datafusion_scalar_udf__"):
+ if hasattr(args[0], "__datafusion_scalar_udf__") or
_is_pycapsule(args[0]):
Review 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`:
```python
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`.
##########
crates/core/src/functions.rs:
##########
@@ -754,16 +808,17 @@ pub fn approx_percentile_cont_with_weight(
}
#[pyfunction]
-#[pyo3(signature = (sort_expression, percentile, filter=None))]
+#[pyo3(signature = (sort_expression, percentile, distinct=None, filter=None))]
pub fn percentile_cont(
sort_expression: PySortExpr,
percentile: f64,
+ distinct: Option<bool>,
filter: Option<PyExpr>,
) -> PyDataFusionResult<PyExpr> {
let agg_fn =
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)
Review 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:
```python
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.
##########
crates/core/src/expr.rs:
##########
@@ -743,6 +767,63 @@ impl PyExpr {
}
}
+/// Start an [`ExprFuncBuilder`] that keeps the options already set on `expr`.
+///
+/// Upstream's `ExprFunctionExt` methods on an `Expr` start from an empty
+/// builder, so `build()` would reset every option not set again. The Python
+/// function wrappers already apply their keyword options, so chaining another
+/// builder method onto their result must not discard them.
+///
+/// A built window function always stores a concrete frame, so whether the user
+/// 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 {
+ match expr {
+ Expr::AggregateFunction(agg) => {
+ let params = &agg.params;
+ let mut builder =
expr.clone().null_treatment(params.null_treatment);
+ if !params.order_by.is_empty() {
+ builder = builder.order_by(params.order_by.clone());
+ }
+ if let Some(filter) = ¶ms.filter {
+ builder = builder.filter(filter.as_ref().clone());
+ }
+ if params.distinct {
+ builder = builder.distinct();
+ }
+ builder
+ }
+ Expr::WindowFunction(window) => {
+ let params = &window.params;
+ let mut builder =
expr.clone().null_treatment(params.null_treatment);
+ if !params.partition_by.is_empty() {
+ builder = builder.partition_by(params.partition_by.clone());
+ }
+ let has_order_by = !params.order_by.is_empty();
+ if has_order_by {
+ builder = builder.order_by(params.order_by.clone());
+ }
+ // A frame equal to the default `build()` derived from the
order-by is
+ // left unset, so it is derived again from the final order-by.
+ if keep_window_frame
+ || params.window_frame
+ !=
datafusion::logical_expr::WindowFrame::new(has_order_by.then_some(true))
Review Comment:
`build()` derives `WindowFrame::new(Some(false))` (RANGE ... CURRENT ROW)
when `order_by` is an empty list, but this check only treats
`WindowFrame::new(None)` as the default when there's no order-by. So a frame
derived from `order_by=[]` looks explicit and survives later chaining:
```python
from datafusion import 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):
return df.select(col("i"),
e.alias("r")).sort(col("i")).collect_column("r").to_pylist()
print(run(f.sum(col("v")).over(Window(order_by=[])).over(Window(order_by="i"))))
# PR: [3, 3, 6, 10] (RANGE frame kept, ties summed); main: [1, 3, 6, 10]
print(run(f.sum(col("v")).over(Window()).over(Window(order_by="i"))))
# PR and main: [1, 3, 6, 10]
```
An empty sort list is easy to get from a dynamically built one. Also
treating `WindowFrame::new(Some(false))` as the default when there's no
order-by would fix it.
##########
docs/source/user-guide/upgrade-guides.md:
##########
@@ -198,6 +198,37 @@ ctx.execute(plan, partitions=0) # before
ctx.execute(plan, partition=0) # after
```
+### More aggregate functions accept `distinct`
Review 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:
```python
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.
##########
python/datafusion/dataframe.py:
##########
@@ -1217,6 +1246,13 @@ def explain(
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.
Review 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:
```python
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.
##########
python/datafusion/dataframe.py:
##########
@@ -1875,6 +1924,34 @@ def fill_null(self, value: Any, subset: list[str] | None
= None) -> DataFrame:
"""
return DataFrame(self.df.fill_null(value, subset))
+ def fill_nan(self, value: float, subset: list[str] | None = None) ->
DataFrame:
+ """Fill NaN values in floating-point columns with a value.
+
+ Only floating-point columns are changed; others are kept unchanged, as
is
+ any column ``value`` cannot be cast to. NaN is distinct from null,
which
+ :py:meth:`fill_null` handles.
+
+ 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
Review 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":
```python
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.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]