namanjain24-sudo commented on code in PR #25091:
URL: https://github.com/apache/datafusion/pull/25091#discussion_r4081974795


##########
datafusion/substrait/src/logical_plan/consumer/rel/set_rel.rs:
##########
@@ -77,16 +81,129 @@ async fn intersect_rels(
     let mut rel = consumer.consume_rel(&rels[0]).await?;
 
     for input in &rels[1..] {
-        rel = LogicalPlanBuilder::intersect(
-            rel,
-            consumer.consume_rel(input).await?,
-            is_all,
-        )?;
+        rel = intersect_rel(rel, consumer.consume_rel(input).await?, is_all)?;
     }
 
     Ok(rel)
 }
 
+/// Intersects two relations, giving the result the nullability the Substrait
+/// [Set Operation rules] prescribe.
+///
+/// [`LogicalPlanBuilder::intersect`] compiles an intersection into a left semi
+/// join, so on its own the result keeps the left input's nullability. The join
+/// matches nulls with nulls, so a left row holding a null in some field only
+/// survives when the right input holds a null there too. A field is therefore
+/// nullable in the result only when it is nullable in *both* inputs.
+///
+/// Applied to each step of a chain, that gives the spec's rule for the 
multiset
+/// intersections - a field is required when any input requires it. For
+/// `INTERSECTION_PRIMARY` the right side is the union of the secondary inputs,
+/// whose field is nullable exactly when some secondary input makes it 
nullable,
+/// so the same rule yields "nullable in the primary input and in at least one
+/// secondary input".
+///
+/// When the right input requires a field the left input leaves nullable, the
+/// intersection is built as an inner join against the distinct right rows
+/// instead, and that field is read from the right side. Matched rows hold 
equal
+/// values, so the result is unchanged, and the field is non-nullable because
+/// its source is: the logical and the physical planner both derive that from
+/// the input schema, so the plan, the physical plan and the batches agree.
+/// Joining against distinct right rows keeps each left row at most once, as 
the
+/// semi join does.
+///
+/// Differing metadata does not change which path is taken, as it is no part of
+/// nullability. The result should describe the left input, exactly as
+/// [`LogicalPlanBuilder::intersect`] does: a column read from the right is 
cast
+/// to an explicit target field carrying the left field's metadata, and an
+/// explicit-field cast target replaces the source's metadata outright rather
+/// than merging into it, so a key only the right input carries is dropped
+/// rather than surviving into the result. A column read from the left is
+/// already exactly the left input's own field, in the logical and in the
+/// physical plan alike, so it needs no such treatment. The physical plan and
+/// the batches carry the same metadata as the logical plan.
+///
+/// [Set Operation rules]: 
https://substrait.io/relations/logical_relations/#set-operation
+fn intersect_rel(
+    left: LogicalPlan,
+    right: LogicalPlan,
+    is_all: bool,
+) -> datafusion::common::Result<LogicalPlan> {
+    let left_fields = left.schema().fields();
+    let right_fields = right.schema().fields();
+    // A field is read from the right side when the left leaves it nullable and
+    // the right requires it. Its metadata does not matter here: it is read 
from
+    // the right with the left field's metadata substituted for its own.
+    let from_right: Vec<bool> = left_fields
+        .iter()
+        .zip(right_fields.iter())
+        .map(|(left, right)| {
+            left.is_nullable()
+                && !right.is_nullable()
+                && left.data_type() == right.data_type()
+        })
+        .collect();
+
+    // `intersect` also reports inputs of different widths.
+    if left_fields.len() != right_fields.len() || !from_right.contains(&true) {
+        return LogicalPlanBuilder::intersect(left, right, is_all);
+    }
+
+    let (left, right, _) = requalify_sides_if_needed(
+        LogicalPlanBuilder::from(left),
+        LogicalPlanBuilder::from(right),
+    )?;
+    let left = if is_all { left } else { left.distinct()? };
+    let right = right.distinct()?.build()?;
+
+    let left_columns = left.schema().columns();
+    let right_columns = right.schema().columns();
+    let exprs = left
+        .schema()
+        .fields()
+        .iter()
+        .zip(&left_columns)
+        .zip(&right_columns)
+        .zip(&from_right)
+        .map(|(((field, left), right), from_right)| {
+            if *from_right {
+                // `alias_qualified_with_metadata` would not do: 
`Expr::Alias`'s
+                // field derivation extends the aliased expression's own
+                // metadata with the alias's, so a key only the right column
+                // carries would survive alongside the left field's metadata.
+                // An explicit-field `Cast` target's metadata is instead used
+                // exactly as given, in both the logical and the physical
+                // plan, so casting to the left field's type (already proven
+                // equal to the right's) and metadata drops the right's own
+                // metadata outright. The qualifier and name still need
+                // `alias_qualified` on top, since a `Cast`'s own field is not
+                // renamed to its target field's name.
+                let target_field = Arc::new(
+                    Field::new(&left.name, field.data_type().clone(), false)
+                        .with_metadata(field.metadata().clone()),
+                );
+                Expr::Cast(Cast::new_from_field(
+                    Box::new(Expr::Column(right.clone())),
+                    target_field,
+                ))
+                .alias_qualified(left.relation.clone(), &left.name)
+            } else {
+                Expr::Column(left.clone())
+            }
+        })
+        .collect::<Vec<_>>();
+
+    left.join_detailed(

Review Comment:
   Thanks — you're right, and I want to flag something more fundamental I found 
while fixing it.
   
   I set the projection's schema explicitly (`join.schema` metadata replaced 
with the left's, same idea as the field-level fix), and it works for the plan 
`intersect_rel` returns. But it doesn't survive `ctx.state().optimize(...)`: 
`LogicalPlan::recompute_schema` rebuilds a `Join`'s or a `Projection`'s 
schema-level metadata from `build_join_schema`/`projection_schema` on its 
current children whenever the optimizer decides *anything* in the plan changed 
— even something unrelated elsewhere in the query. I confirmed this 
reintroduces `only_in_secondary` even for the simplest two-input case, no union 
or join filter involved, so it isn't specific to this construction. Field-level 
metadata is unaffected, since it's carried on the `Field` objects themselves 
and survives any rebuild — that part is solid.
   
   I don't see a supported way from this crate to opt a join's or a 
projection's schema-level metadata out of that recomputation. Pushed 3c11ef9:
   
   - Documented this in `intersect_rel`'s doc comment.
   - Extended the test to check the complete metadata map (not one key), so it 
now catches `only_in_secondary` — on the plan right after conversion, which is 
what this crate controls.
   - Since physical planning runs on the *optimized* plan, changed the 
physical/batch comparison to check against the optimized logical schema instead 
of the pre-optimization one. Field-level metadata and nullability are still 
asserted at every stage (conversion, optimized, physical, batches) and catch a 
regression in either.
   
   So schema-level metadata is correct right after `from_substrait_plan`, but a 
caller that runs the plan through the optimizer may still see 
`only_in_secondary` on the final schema. If that's not acceptable, I think it 
needs either a DataFusion-side change (a way to pin schema-level metadata 
through `recompute_schema`) or a different construction than an inner join for 
the narrowing path — both bigger than this PR. Let me know how you'd like to 
proceed.



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