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]

Reply via email to