alexy opened a new issue, #26065:
URL: https://github.com/apache/datafusion/issues/26065
### Is your feature request related to a problem or challenge?
Physical planning gets slow, and grows faster than the query does, when
projections carry many typed NULL struct literals. A typical source is a `UNION
ALL` that packs rows of different kinds into one relation, one struct column
per kind, where each branch fills its own column and puts `CAST(NULL AS
STRUCT<...>)` in the others.
`ProjectionExec::try_from_projector` calls `EquivalenceProperties::project`,
which registers every literal in the projection as a constant. For each new
uniform constant, `EquivalenceGroup::add_constant` compares its value with the
constant of every existing class
(`datafusion/physical-expr/src/equivalence/class.rs`, the loop over
`self.classes`). For nested values, `ScalarValue::eq` compares the inner arrays
with arrow's `PartialEq`, which is `self.to_data() == other.to_data()`: both
arrays are converted to `ArrayData`, including every child, before anything is
compared, even when the two values have different types. A projection with n
struct literals pays for n² of these conversions, and planners call
`try_from_projector` again whenever a projection's children are replaced
(`EnsureRequirements`, filter pushdown, projection pushdown).
### To reproduce
On `main` (8248a57969), with branch i filling column i with a struct and
every other column with a typed NULL struct, one struct type per column:
```rust
use std::time::Instant;
use datafusion::prelude::*;
fn struct_type(col: usize, fields: usize) -> String {
let f: Vec<String> = (0..fields).map(|j| format!("c{col}_f{j}
INT")).collect();
format!("STRUCT<{}>", f.join(", "))
}
fn query(cols: usize, fields: usize) -> String {
(0..cols)
.map(|i| {
let exprs: Vec<String> = (0..cols)
.map(|k| if k == i {
let f: Vec<String> = (0..fields).map(|j|
format!("'c{k}_f{j}', {j}")).collect();
format!("CAST(named_struct({}) AS {}) AS s{k}",
f.join(", "), struct_type(k, fields))
} else {
format!("CAST(NULL AS {}) AS s{k}", struct_type(k,
fields))
})
.collect();
format!("SELECT {i} AS kind, {} FROM t", exprs.join(", "))
})
.collect::<Vec<_>>()
.join("\nUNION ALL\n")
}
#[tokio::main]
async fn main() -> datafusion::error::Result<()> {
let ctx = SessionContext::new();
ctx.sql("CREATE TABLE t AS VALUES (1)").await?.collect().await?;
for (cols, fields) in [(10, 20), (20, 20), (35, 20), (35, 40)] {
let df = ctx.sql(&query(cols, fields)).await?;
let started = Instant::now();
df.create_physical_plan().await?;
println!("{cols} columns x {fields} fields: {:?}",
started.elapsed());
}
Ok(())
}
```
`create_physical_plan` alone, release build, Apple M1 Max:
| struct columns x fields | main | with the change below |
|---|---:|---:|
| 10 x 20 | 31 ms | 17 ms |
| 20 x 20 | 179 ms | 47 ms |
| 35 x 20 | 862 ms | 147 ms |
| 35 x 40 | 1,766 ms | 253 ms |
We found it through a real workload: a game's world kept as one relation of
35 kinds, planned per step through [Sail](https://github.com/lakehq/sail) on
DataFusion 55.1. A profile of that plan put about a quarter of all planning
samples under `Literal::dyn_eq` -> `StructArray::eq` ->
`ArrayData::from(StructArray)`, reached from `EquivalenceProperties::project`
and `EquivalenceGroup::add_constant`.
### Describe the solution you'd like
Compare cheaply before comparing data. In `ScalarValue::eq`, for the nested
variants (`List`, `LargeList`, `FixedSizeList`, `ListView`, `LargeListView`,
`Struct`, `Map`):
```rust
fn nested_eq<T: Array + PartialEq>(a: &Arc<T>, b: &Arc<T>) -> bool {
Arc::ptr_eq(a, b)
|| (a.len() == b.len() && a.data_type() == b.data_type() &&
a.as_ref() == b.as_ref())
}
```
This gives the numbers in the right-hand column, and `cargo test -p
datafusion-common --lib` passes (623 tests). Values of the same type still
convert to `ArrayData`; `add_constant` could also avoid the linear scan, for
example by keying uniform constants by value, but the type check alone removes
most of the cost here.
### Describe alternatives you've considered
Rewriting the SQL. Placeholders that cannot be folded into literals (`CASE
WHEN <never> THEN named_struct(<typed NULLs>) END`) avoid the constant handling
but make the plan larger, and planned slower in our workload. arrow-rs could
also check `data_type()` first in its array `PartialEq`; that would help every
caller.
### Additional context
Part of the planning-speed work in #19795. Happy to open a PR with the
change and a test.
--
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]