namanjain24-sudo opened a new issue, #25366:
URL: https://github.com/apache/datafusion/issues/25366

   ### Describe the bug
   
   The Substrait producer leaves `output_type` unset on two kinds of expression 
it emits:
   
   - window functions (`producer/expr/window_function.rs:110`)
   - `LIKE` / `ILIKE`, and the `not` wrapper around a negated one 
(`producer/expr/scalar_function.rs:296` and `:312`)
   
   Both `Expression.ScalarFunction.output_type` and 
`Expression.WindowFunction.output_type` are documented as:
   
   > Must be set to the return type of the function, exactly as derived using 
the declaration in the extension.
   
   The wording is the same in the `substrait` 0.63.0 crate this repo pins and 
in substrait-io/substrait `main`, so this is not a version difference.
   
   substrait-java reads both fields with the same converter that rejects an 
unset type, so it refuses these plans.
   
   Reproduced on `main` at 86ba5e0.
   
   ### To Reproduce
   
   Add this as an example under `datafusion/substrait/examples/` and run
   `cargo run --locked -p datafusion-substrait --example output_type_probe`. It 
inspects the produced protobuf directly, without converting it back through a 
consumer.
   
   ```rust
   use datafusion::arrow::datatypes::{DataType, Field, Schema};
   use datafusion::common::Result;
   use datafusion::datasource::empty::EmptyTable;
   use datafusion::prelude::SessionContext;
   use datafusion_substrait::logical_plan::producer::to_substrait_plan;
   use datafusion_substrait::substrait::proto::expression::RexType;
   use datafusion_substrait::substrait::proto::rel::RelType;
   use datafusion_substrait::substrait::proto::{plan_rel, Rel};
   use std::sync::Arc;
   
   fn walk(rel: &Rel, out: &mut Vec<String>) {
       match rel.rel_type.as_ref() {
           Some(RelType::Project(p)) => {
               for e in &p.expressions {
                   match e.rex_type.as_ref() {
                       Some(RexType::WindowFunction(w)) => out.push(format!(
                           "WindowFunction.output_type set = {}",
                           w.output_type.is_some()
                       )),
                       Some(RexType::ScalarFunction(f)) => out.push(format!(
                           "ScalarFunction.output_type  set = {}",
                           f.output_type.is_some()
                       )),
                       _ => {}
                   }
               }
               p.input.as_ref().map(|i| walk(i, out));
           }
           Some(RelType::Filter(f)) => {
               if let Some(RexType::ScalarFunction(s)) =
                   f.condition.as_ref().and_then(|c| c.rex_type.as_ref())
               {
                   out.push(format!(
                       "ScalarFunction.output_type  set = {}",
                       s.output_type.is_some()
                   ));
               }
               f.input.as_ref().map(|i| walk(i, out));
           }
           Some(RelType::Aggregate(a)) => {
               a.input.as_ref().map(|i| walk(i, out));
           }
           _ => {}
       };
   }
   
   #[tokio::main(flavor = "current_thread")]
   async fn main() -> Result<()> {
       let ctx = SessionContext::new();
       ctx.register_table(
           "t",
           Arc::new(EmptyTable::new(Arc::new(Schema::new(vec![Field::new(
               "i",
               DataType::Int64,
               true,
           )])))),
       )?;
   
       for sql in [
           "SELECT sum(i) OVER (ORDER BY i) FROM t",
           "SELECT i FROM t WHERE CAST(i AS VARCHAR) LIKE '1%'",
           "SELECT i FROM t WHERE CAST(i AS VARCHAR) NOT LIKE '1%'",
           "SELECT i + 1 FROM t",
       ] {
           let plan = ctx.sql(sql).await?.into_optimized_plan()?;
           let proto = to_substrait_plan(&plan, &ctx.state())?;
           let mut out = vec![];
           for r in &proto.relations {
               if let Some(plan_rel::RelType::Root(root)) = &r.rel_type {
                   root.input.as_ref().map(|i| walk(i, &mut out));
               }
           }
           for line in out {
               println!("{sql:<52} {line}");
           }
       }
       Ok(())
   }
   ```
   
   Output:
   
   ```
   SELECT sum(i) OVER (ORDER BY i) FROM t               
WindowFunction.output_type set = false
   SELECT i FROM t WHERE CAST(i AS VARCHAR) LIKE '1%'   
ScalarFunction.output_type  set = false
   SELECT i FROM t WHERE CAST(i AS VARCHAR) NOT LIKE '1%' 
ScalarFunction.output_type  set = false
   SELECT i + 1 FROM t                                  
ScalarFunction.output_type  set = true
   ```
   
   The last line is the comparison: #20597 set the type for binary expressions, 
which is why `i + 1` carries one.
   
   ### Expected behavior
   
   `Expression.WindowFunction.output_type` and the `LIKE` / `ILIKE` 
`Expression.ScalarFunction.output_type` should be set to the expression's 
return type, the way `from_binary_expr` and the scalar function helpers already 
do.
   
   ### Additional context
   
   **What breaks.** substrait-java 0.103.0 rejects both plans. Taking the plans 
this repo produces for the two queries above and feeding them to 
`ProtoPlanConverter` (with the `extension_urn_reference` patched first, 
otherwise #11545 stops it earlier):
   
   | query | error |
   | --- | --- |
   | `SELECT sum(i) OVER (ORDER BY i) FROM t` | `UnsupportedOperationException: 
Type is not set` at `ProtoExpressionConverter.fromWindowFunction:406` |
   | `... WHERE CAST(i AS VARCHAR) LIKE '1%'` | `UnsupportedOperationException: 
Type is not set` at `ProtoExpressionConverter:201`, via 
`ProtoRelConverter.newFilter` |
   
   Both come from `ProtoTypeConverter.from:125`, the same code path that 
rejects an aggregate with no `output_type` in #25049.
   
   **Why no test catches it.** The DataFusion consumer reads `output_type` only 
in `consumer/expr/cast.rs`; nothing in the window function or scalar function 
path reads it. Producer and consumer therefore agree on the omission and every 
DataFusion-to-DataFusion round trip passes.
   
   **Related, but distinct:**
   
   - #25049 and #25090 are the same field on `AggregateFunction`, and neither 
touches the window function or `LIKE` sites.
   - #15831 and #20597 set this field for binary and unary expressions, 
`from_function`, and higher order functions. #20597 edited this same file but 
left the `LIKE` helper and the window function as they were.
   - #25100 and #25146 are about `phase` on the same `WindowFunction` message, 
not `output_type`.
   
   One more site I did not include: `producer/utils.rs:106`, the shared 
`negate()` helper, has the same `output_type: None`. Its callers are the 
negated subquery forms, and the optimizer decorrelates `NOT EXISTS` and `NOT IN 
(<subquery>)` into anti joins before the producer sees them, so I could not 
produce a plan from SQL that reaches it. `NOT IN (1, 2)`, `NOT BETWEEN` and `IS 
NOT NULL` are rewritten into `and`/`not_equal`, `or`/`lt`/`gt` and 
`is_not_null`, which all carry a type already.


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