Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
31 commits
Select commit Hold shift + click to select a range
5a0e79d
feat: expose any_value, array math, array_first, and missing Spark fu…
timsaucer Sep 25, 2026
8027fcb
feat: add rand, substring_index, file metadata functions, and distinc…
timsaucer Sep 25, 2026
97e6031
feat(spark): add pyspark aliases and optional-length substr
timsaucer Sep 25, 2026
e8c2192
feat: accept optional arguments upstream supports on trim, array_to_s…
timsaucer Sep 25, 2026
2747573
feat: add DataFrame.fill_nan
timsaucer Sep 25, 2026
24764bd
feat: add null_treatment to lead and lag
timsaucer Sep 25, 2026
b035c97
feat: add show_statistics, analyze_level, and analyze_categories to e…
timsaucer Sep 25, 2026
d2ac890
fix: keep existing options when chaining aggregate and window builders
timsaucer Sep 25, 2026
97e5905
feat: accept a bare PyCapsule in ScalarUDF and WindowUDF
timsaucer Sep 25, 2026
d53a255
chore: export TableProviderFactory from datafusion.catalog
timsaucer Sep 25, 2026
5d43576
docs(skill): record DataFrame.to_string as not needing exposure
timsaucer Sep 25, 2026
f6d0478
fix: resolve the _PyCapsule type alias so capsule overloads type-check
timsaucer Sep 25, 2026
6345b17
fix: keep an explicit window frame that equals the default when chaining
timsaucer Sep 25, 2026
73b052d
fix: leave lead and lag null_treatment unset by default
timsaucer Sep 25, 2026
8173554
fix: keep existing window function options in over()
timsaucer Sep 25, 2026
96f990b
feat: accept native ints and a single argument in range and gen_series
timsaucer Sep 25, 2026
d731279
fix: decide window frame re-derivation from the expression alone
timsaucer Sep 28, 2026
b4cad2d
fix: treat the frame an empty order_by derives as a default when chai…
timsaucer Sep 28, 2026
92ace16
fix: raise when a builder option does not apply to the function kind
timsaucer Sep 28, 2026
8e21c8a
docs: note that over() drops options set on an aggregate
timsaucer Sep 28, 2026
6b52225
docs: document that chaining now keeps options already set
timsaucer Sep 28, 2026
0ac494e
fix: reject a non-bool distinct in string_agg
timsaucer Sep 28, 2026
607dde4
fix: reject show_statistics combined with analyze in explain
timsaucer Sep 28, 2026
4933fa2
fix: treat a bare str as a column name in spark.printf
timsaucer Sep 28, 2026
23ed533
fix: name the expected and found capsule when importing a UDF
timsaucer Sep 28, 2026
8454a09
fix: treat an empty fill_null or fill_nan subset as no columns
timsaucer Sep 28, 2026
1b01845
fix: keep the sort direction in percentile_cont and friends
timsaucer Sep 28, 2026
f8d7b70
fix: allow the keyword-only form of the udf, udaf, and udwf decorators
timsaucer Sep 28, 2026
7259819
docs: fix the broken guide links in aggregate docstrings
timsaucer Sep 28, 2026
d779502
test: check that bit_and and bit_or keep distinct in the expression
timsaucer Sep 28, 2026
db5b7cc
docs(skill): show substr's length argument as valid
timsaucer Sep 28, 2026
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions .ai/skills/check-upstream/SKILL.md
Original file line number Diff line number Diff line change
Expand Up @@ -146,6 +146,7 @@ The user may specify an area via `$ARGUMENTS`. If no area is specified or "all"
- `show_limit` — already covered by `DataFrame.show()`, which provides the same functionality with a simpler API
- `with_param_values` — already covered by the `param_values` argument on `SessionContext.sql()`, which accomplishes the same thing more robustly
- `union_by_name_distinct` — already covered by `DataFrame.union_by_name(distinct=True)`, which provides a more Pythonic API
- `to_string` — `str(df)` is the Pythonic way to get a string and already goes through `__repr__` and the configurable formatter. A separate `to_string()` would either duplicate `str(df)` or render every row through a different path (session `datafusion.format.*` options, no formatter), giving a third text rendering alongside `repr` and `show()`

**How to check:**
1. Fetch the upstream DataFrame documentation page listing all methods
Expand Down
49 changes: 47 additions & 2 deletions crates/core/src/dataframe.rs
Original file line number Diff line number Diff line change
Expand Up @@ -844,13 +844,24 @@ impl PyDataFrame {
}

/// Print the query plan
#[pyo3(signature = (verbose=false, analyze=false, format=None))]
#[pyo3(signature = (
verbose=false,
analyze=false,
format=None,
show_statistics=None,
analyze_level=None,
analyze_categories=None
))]
#[allow(clippy::too_many_arguments)]
fn explain(
&self,
py: Python,
verbose: bool,
analyze: bool,
format: Option<&str>,
show_statistics: Option<bool>,
analyze_level: Option<&str>,
analyze_categories: Option<Vec<String>>,
) -> PyDataFusionResult<()> {
let explain_format = match format {
Some(f) => f
Expand All @@ -860,10 +871,24 @@ impl PyDataFrame {
})?,
None => datafusion::common::format::ExplainFormat::Indent,
};
let analyze_level = analyze_level
.map(|l| l.parse::<datafusion::common::format::MetricType>())
.transpose()?;
let analyze_categories = analyze_categories
.map(|cats| {
cats.iter()
.map(|c| c.parse::<datafusion::common::format::MetricCategory>())
.collect::<datafusion::common::Result<Vec<_>>>()
.map(datafusion::common::format::ExplainAnalyzeCategories::Only)
})
.transpose()?;
let opts = datafusion::logical_expr::ExplainOption::default()
.with_verbose(verbose)
.with_analyze(analyze)
.with_format(explain_format);
.with_format(explain_format)
.with_show_statistics(show_statistics)
.with_analyze_level(analyze_level)
.with_analyze_categories(analyze_categories);
let df = self.df.as_ref().clone().explain_with_options(opts)?;
print_dataframe(py, df)
}
Expand Down Expand Up @@ -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)?;

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.

Ok(Self::new(df))
}
}

#[derive(Debug, Clone, PartialEq, Eq, Hash, PartialOrd, Ord)]
Expand Down
112 changes: 101 additions & 11 deletions crates/core/src/expr.rs
Original file line number Diff line number Diff line change
Expand Up @@ -55,7 +55,7 @@ use crate::expr::aggregate_expr::PyAggregateFunction;
use crate::expr::binary_expr::PyBinaryExpr;
use crate::expr::column::PyColumn;
use crate::expr::literal::PyLiteral;
use crate::functions::add_builder_fns_to_window;
use crate::functions::{add_builder_fns_to_window, apply_window_options};
use crate::pyarrow_util::scalar_to_pyarrow;
use crate::sql::logical::PyLogicalPlan;

Expand Down Expand Up @@ -625,34 +625,61 @@ impl PyExpr {
// Expression Function Builder functions

pub fn order_by(&self, order_by: Vec<PySortExpr>) -> PyExprFuncBuilder {
self.expr
.clone()
builder_from_expr(&self.expr)
.order_by(to_sort_expressions(order_by))
.into()
}

pub fn filter(&self, filter: PyExpr) -> PyExprFuncBuilder {
self.expr.clone().filter(filter.expr.clone()).into()
builder_from_expr(&self.expr)
.filter(filter.expr.clone())
.into()
}

pub fn distinct(&self) -> PyExprFuncBuilder {
self.expr.clone().distinct().into()
// Only aggregates support DISTINCT, including an aggregate run as a window
// function. For anything else, upstream's empty builder makes `build()`
// raise instead of dropping the option.
let supports_distinct = match &self.expr {
Expr::AggregateFunction(_) => true,
Expr::WindowFunction(window) => {
matches!(window.fun, WindowFunctionDefinition::AggregateUDF(_))
}
_ => false,
};
if supports_distinct {
builder_from_expr(&self.expr).distinct().into()
} else {
self.expr.clone().distinct().into()
}
}

pub fn null_treatment(&self, null_treatment: NullTreatment) -> PyExprFuncBuilder {
self.expr
.clone()
builder_from_expr(&self.expr)
.null_treatment(Some(null_treatment.into()))
.into()
}

pub fn partition_by(&self, partition_by: Vec<PyExpr>) -> PyExprFuncBuilder {
let partition_by = partition_by.iter().map(|e| e.expr.clone()).collect();
self.expr.clone().partition_by(partition_by).into()
// Window-only option: on an aggregate, upstream's empty builder makes
// `build()` raise instead of dropping it.
if matches!(self.expr, Expr::AggregateFunction(_)) {
return self.expr.clone().partition_by(partition_by).into();
}
builder_from_expr(&self.expr)
.partition_by(partition_by)
.into()
}

pub fn window_frame(&self, window_frame: PyWindowFrame) -> PyExprFuncBuilder {
self.expr.clone().window_frame(window_frame.into()).into()
// Window-only option; see `partition_by`.
if matches!(self.expr, Expr::AggregateFunction(_)) {
return self.expr.clone().window_frame(window_frame.into()).into();
}
builder_from_expr(&self.expr)
.window_frame(window_frame.into())
.into()
}

#[pyo3(signature = (partition_by=None, window_frame=None, order_by=None, null_treatment=None))]
Expand All @@ -678,8 +705,8 @@ impl PyExpr {
null_treatment,
)
}
Expr::WindowFunction(_) => add_builder_fns_to_window(
self.expr.clone(),
Expr::WindowFunction(_) => apply_window_options(
builder_from_expr(&self.expr),
partition_by,
window_frame,
order_by,
Expand Down Expand Up @@ -743,6 +770,69 @@ 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. A frame equal to the default for the current order-by is
/// treated as unset, which depends only on the expression and so behaves the
/// same after a copy, pickle, or round trip through protobuf or SQL.
fn builder_from_expr(expr: &Expr) -> 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) = &params.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. 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 {
params.window_frame == datafusion::logical_expr::WindowFrame::new(Some(true))
} else {
params.window_frame == datafusion::logical_expr::WindowFrame::new(None)
|| params.window_frame
== datafusion::logical_expr::WindowFrame::new(Some(false))
};
if !is_default_frame {
builder = builder.window_frame(params.window_frame.clone());
}
if let Some(filter) = &params.filter {
builder = builder.filter(filter.as_ref().clone());
}
if params.distinct {
builder = builder.distinct();
}
builder
}
_ => expr.clone().null_treatment(None),
}
}

#[pyclass(
from_py_object,
frozen,
Expand Down
Loading
Loading