jayzhan211 commented on code in PR #25391:
URL: https://github.com/apache/datafusion/pull/25391#discussion_r4225946992


##########
datafusion/optimizer/src/decorrelate_predicate_subquery.rs:
##########
@@ -783,13 +784,16 @@ fn build_join(
 ) -> Result<Option<BuiltJoin>> {
     let mut pull_up = PullUpCorrelatedExpr::new()
         .with_in_predicate_opt(in_predicate_opt.cloned())
-        .with_exists_sub_query(in_predicate_opt.is_none());
+        .with_exists_sub_query(in_predicate_opt.is_none())
+        .with_need_handle_count_bug(true);

Review Comment:
   `with_need_handle_count_bug(true)` also applies to `IN`/`NOT IN`, so every 
correlated `IN` over a groupless aggregate now hits the `not_impl_err` at 
`:1107`. That includes `max`/`sum`, which main answers correctly.
   
   ```sql
   -- main and DuckDB: 11; this PR: "IN/NOT IN count-bug compensation ... is 
not implemented"
   SELECT t1_id FROM t1
   WHERE t1_int + 2 IN (SELECT max(t2_int) FROM t2 WHERE t2.t2_id = t1.t1_id);
   -- main and DuckDB: 22, 33, 44; this PR: same error
   SELECT t1_id FROM t1
   WHERE t1_int NOT IN (SELECT count(*) FROM t2 WHERE t2.t2_id = t1.t1_id);
   ```
   
   Keep IN on the old path until the follow-up, drop the `not_impl_err` branch, 
and add both queries to `subquery.slt`



##########
datafusion/optimizer/src/decorrelate_predicate_subquery.rs:
##########
@@ -1063,6 +1089,116 @@ fn build_join(
     }))
 }
 
+/// Builds the join for a correlated `EXISTS` subquery whose groupless
+/// aggregate requires count-bug compensation.
+#[expect(clippy::too_many_arguments)]
+fn build_join_with_count_bug(
+    left: &LogicalPlan,
+    sub_query_alias: LogicalPlan,
+    join_filter_opt: Option<Expr>,
+    in_predicate_opt: Option<&Expr>,
+    join_type: JoinType,
+    alias: &str,
+    expr_map: &crate::decorrelate::ExprResultMap,
+    pull_up_having_expr: Option<&Expr>,
+    forces_empty_result: bool,
+) -> Result<LogicalPlan> {
+    if in_predicate_opt.is_some() {
+        return not_impl_err!(
+            "IN/NOT IN count-bug compensation for groupless aggregates is not 
implemented"
+        );
+    }
+
+    let joined = LogicalPlanBuilder::from(left.clone())
+        .join_on(sub_query_alias, JoinType::Left, join_filter_opt)?
+        .build()?;
+
+    let left_projection: Vec<Expr> = left
+        .schema()
+        .columns()
+        .into_iter()
+        .map(Expr::from)
+        .collect();
+
+    let indicator_col = Expr::Column(Column::new(Some(alias), 
UN_MATCHED_ROW_INDICATOR));
+
+    let having_arm = pull_up_having_expr.map(|f| f.clone().is_not_true());
+
+    let mut expr_rewrite = TypeCoercionRewriter::new(joined.schema());
+
+    // The subquery's HAVING clause, evaluated against an unmatched row's
+    // default aggregate values (e.g. count(*) defaults to 0, sum(x) to
+    // NULL), is usually `true`, but not once a second filter is combined
+    // into it, so it cannot be assumed.
+    let unmatched_having_default = match pull_up_having_expr {
+        Some(f) => f
+            .clone()
+            .transform_up(|e| {
+                if let Expr::Column(Column { name, .. }) = &e
+                    && let Some(default_value) = expr_map.get(name)
+                {
+                    return Ok(Transformed::yes(default_value.clone()));
+                }
+                Ok(Transformed::no(e))
+            })
+            .data()
+            .and_then(|simplified| simplified.rewrite(&mut 
expr_rewrite).data())?,
+        None => lit(true),
+    };
+
+    // EXISTS is true by default (the groupless aggregate always
+    // produces a row), unless one of three cases holds:
+    // (1) the outer row joined against an inner group whose HAVING
+    //     predicate failed,
+    // (2) an unmatched outer row's own default aggregate values fail
+    //     that same HAVING clause,
+    // (3) the subquery was truncated to zero rows unconditionally
+    //     (LIMIT 0), which is false for every outer row regardless of
+    //     whether it matched.
+    let exists_expr = if forces_empty_result {
+        lit(false)
+    } else {
+        match having_arm {
+            Some(when_expr) => {
+                when(indicator_col.clone().is_null(), unmatched_having_default)

Review Comment:
   The unmatched arm returns the HAVING evaluated on the defaults as-is, which 
is NULL when it reads `sum`/`max`. `LeftAnti` keeps a row only when `NOT 
exists_expr` is true, so a NULL drops it; a mark join's `NOT mark` is NULL too. 
NOT EXISTS loses rows whose subquery is empty.
   
   ```sql
   -- DuckDB and main: 11, 22, 33, 44; this PR: 11, 22, 44 (same with `t1_id > 
100 OR NOT EXISTS`)
   SELECT t1_id FROM t1 WHERE NOT EXISTS (
     SELECT c FROM (
       SELECT count(*) AS c, sum(t2_int) AS s FROM t2
       WHERE t1.t1_id = t2.t2_id HAVING count(*) = 0
     ) agg WHERE s > 5
   );
   ```



##########
datafusion/optimizer/src/decorrelate.rs:
##########
@@ -351,6 +368,49 @@ impl TreeNodeRewriter for PullUpCorrelatedExpr {
                     None
                 };
 
+                if self.pull_up_having_expr.is_some()

Review Comment:
   `ScalarSubqueryToJoin` also runs this branch, and it returns the empty-input 
default for unmatched rows without checking `pull_up_having_expr`. When a 
stacked filter that the default fails is merged in, the scalar subquery returns 
the default instead of NULL.
   
   ```sql
   -- DuckDB and main: v is NULL for every row; this PR: v = 0 for t1_id = 33
   SELECT t1_id, (
     SELECT c FROM (
       SELECT count(*) AS c FROM t2 WHERE t2.t2_id = t1.t1_id HAVING count(*) < 
10
     ) s WHERE c > 5
   ) AS v FROM t1;
   ```



##########
datafusion/optimizer/src/decorrelate_predicate_subquery.rs:
##########
@@ -1063,6 +1089,116 @@ fn build_join(
     }))
 }
 
+/// Builds the join for a correlated `EXISTS` subquery whose groupless
+/// aggregate requires count-bug compensation.
+#[expect(clippy::too_many_arguments)]
+fn build_join_with_count_bug(
+    left: &LogicalPlan,
+    sub_query_alias: LogicalPlan,
+    join_filter_opt: Option<Expr>,
+    in_predicate_opt: Option<&Expr>,
+    join_type: JoinType,
+    alias: &str,
+    expr_map: &crate::decorrelate::ExprResultMap,
+    pull_up_having_expr: Option<&Expr>,
+    forces_empty_result: bool,
+) -> Result<LogicalPlan> {
+    if in_predicate_opt.is_some() {
+        return not_impl_err!(
+            "IN/NOT IN count-bug compensation for groupless aggregates is not 
implemented"
+        );
+    }
+
+    let joined = LogicalPlanBuilder::from(left.clone())
+        .join_on(sub_query_alias, JoinType::Left, join_filter_opt)?
+        .build()?;
+
+    let left_projection: Vec<Expr> = left
+        .schema()
+        .columns()
+        .into_iter()
+        .map(Expr::from)
+        .collect();
+
+    let indicator_col = Expr::Column(Column::new(Some(alias), 
UN_MATCHED_ROW_INDICATOR));
+
+    let having_arm = pull_up_having_expr.map(|f| f.clone().is_not_true());
+
+    let mut expr_rewrite = TypeCoercionRewriter::new(joined.schema());
+
+    // The subquery's HAVING clause, evaluated against an unmatched row's
+    // default aggregate values (e.g. count(*) defaults to 0, sum(x) to
+    // NULL), is usually `true`, but not once a second filter is combined
+    // into it, so it cannot be assumed.
+    let unmatched_having_default = match pull_up_having_expr {
+        Some(f) => f
+            .clone()
+            .transform_up(|e| {
+                if let Expr::Column(Column { name, .. }) = &e
+                    && let Some(default_value) = expr_map.get(name)
+                {
+                    return Ok(Transformed::yes(default_value.clone()));
+                }
+                Ok(Transformed::no(e))
+            })
+            .data()
+            .and_then(|simplified| simplified.rewrite(&mut 
expr_rewrite).data())?,
+        None => lit(true),
+    };
+
+    // EXISTS is true by default (the groupless aggregate always
+    // produces a row), unless one of three cases holds:
+    // (1) the outer row joined against an inner group whose HAVING
+    //     predicate failed,
+    // (2) an unmatched outer row's own default aggregate values fail
+    //     that same HAVING clause,
+    // (3) the subquery was truncated to zero rows unconditionally
+    //     (LIMIT 0), which is false for every outer row regardless of
+    //     whether it matched.
+    let exists_expr = if forces_empty_result {
+        lit(false)
+    } else {
+        match having_arm {
+            Some(when_expr) => {
+                when(indicator_col.clone().is_null(), unmatched_having_default)
+                    .when(when_expr, lit(false))
+                    .otherwise(lit(true))?
+                    .rewrite(&mut expr_rewrite)
+                    .data()?
+            }
+            None => lit(true),

Review Comment:
   A HAVING that reads an outer column becomes a join conjunct, so a group that 
exists but fails it is also a LEFT JOIN miss. The `CASE` treats every miss as 
"no group", and with no HAVING left, `exists_expr` is `lit(true)`.
   
   ```sql
   -- DuckDB and main: 11; this PR: 11, 22, 33, 44 (NOT EXISTS returns nothing)
   SELECT t1_id FROM t1 WHERE EXISTS (
     SELECT count(*) FROM t2 WHERE t2.t2_id = t1.t1_id HAVING count(*) = 
t1.t1_int
   );
   ```
   
   Add this to the `compensable` check in `build_join` (with the fallback 
suggested at `decorrelate.rs:419`):
   ```rs
   // A miss is only "no group" if the join compares nothing but group keys.
   let reads_agg_output = 
count_bug_compensation.as_ref().is_some_and(|expr_map| {
       pull_up
           .join_filters
           .iter()
           .flat_map(|f| f.column_refs())
           .any(|c| expr_map.contains_key(&c.name))
   });
   ```



##########
datafusion/optimizer/src/decorrelate.rs:
##########
@@ -677,31 +749,80 @@ impl PullUpCorrelatedExpr {
     }
 
     fn collect_missing_exprs(
-        &self,
+        &mut self,
         exprs: &[Expr],
         correlated_subquery_cols: &BTreeSet<Column>,
+        subquery_schema: &DFSchemaRef,
     ) -> Result<Vec<Expr>> {
         let mut missing_exprs = vec![];
         for expr in exprs {
             if !missing_exprs.contains(expr) {
                 missing_exprs.push(expr.clone())
             }
         }
+        // A correlated column compared bare (`t1.a = t2.b`) only ever
+        // matches one row per value, so grouping the subquery's own
+        // aggregate by that column is enough. But if the join filter
+        // wraps it in an expression (`t1.a = CAST(t2.b AS INT)`), two
+        // different column values can compare equal after the cast, and
+        // grouping by the bare column would compute the aggregate once
+        // per underlying value instead of once per outer row it actually
+        // joins against. Grouping by the wrapping expression instead
+        // keeps those together, matching what the join predicate itself
+        // treats as equal.
+        let join_filter_exprs =
+            collect_subquery_join_exprs(&self.join_filters, subquery_schema)?;
         for col in correlated_subquery_cols.iter() {
-            let col_expr = Expr::Column(col.clone());
-            if !missing_exprs.contains(&col_expr) {
+            if let Some(existing_alias) = self.correlated_col_aliases.get(col) 
{
+                if !collides_with_existing(&missing_exprs, existing_alias) {
+                    missing_exprs.push(existing_alias.clone())
+                }
+                continue;
+            }
+            let wrapped = join_filter_exprs

Review Comment:
   The wrapped-key alias runs in every pull-up (plain EXISTS, IN, scalar), not 
just the compensated aggregate. Once the CAST is aliased, the bare column is no 
longer projected even if another conjunct still reads it. The alias is also 
unqualified and named only after `col.name`, so two scalar subqueries on the 
same column collide.
   
   ```sql
   CREATE TABLE o5(a BIGINT, c INT) AS VALUES (1, 0), (2, 5);
   CREATE TABLE i5(b INT) AS VALUES (1), (2);
   -- main: 1; this PR: Schema error: No field named __correlated_sq_1.b
   SELECT a FROM o5 WHERE EXISTS (SELECT 1 FROM i5 WHERE o5.a = i5.b AND o5.c < 
i5.b);
   
   CREATE TABLE o6(x BIGINT) AS VALUES (1), (1);
   CREATE TABLE i6(id INT, k INT) AS VALUES (1, 1);
   -- main: (1, 1, 1) twice; this PR: Ambiguous reference to unqualified field 
__correlated_group_expr_k
   SELECT x,
     (SELECT count(*) FROM i6 WHERE i6.k = o6.x),
     (SELECT sum(id) FROM i6 WHERE i6.k = o6.x)
   FROM o6;
   ```
   
   Suggest dropping the alias rewrite (`collect_subquery_join_exprs`, 
`correlated_col_aliases`, `collides_with_existing`) and falling back to the old 
rewrite when a correlated key is wrapped. The CAST-dedup cases then keep main's 
answer, which can be a follow-up:
   ```rs
   /// Whether every subquery-side operand of `join_filters` is a bare column, 
so a
   /// LEFT JOIN against the subquery grouped by those columns matches at most 
one
   /// group per outer row.
   fn subquery_operands_are_bare_columns(
       join_filters: &[Expr],
       subquery_cols: &BTreeSet<Column>,
   ) -> bool {
       join_filters.iter().all(|filter| {
           let mut bare = true;
           filter
               .apply(|e| {
                   let refs = e.column_refs();
                   if !refs.is_empty() && refs.iter().all(|c| 
subquery_cols.contains(*c)) {
                       bare = matches!(e, Expr::Column(_));
                       return Ok(if bare {
                           TreeNodeRecursion::Jump
                       } else {
                           TreeNodeRecursion::Stop
                       });
                   }
                   Ok(TreeNodeRecursion::Continue)
               })
               .expect("infallible");
           bare
       })
   }
   ```
   Then in `build_join`, require it whenever a count map is present:
   ```rs
   let subquery_cols: BTreeSet<Column> = pull_up
       .correlated_subquery_cols_map
       .values()
       .flatten()
       .cloned()
       .collect();
   let keys_are_bare = count_bug_compensation.is_none()
       || subquery_operands_are_bare_columns(&pull_up.join_filters, 
&subquery_cols);
   ```



##########
datafusion/optimizer/src/decorrelate_predicate_subquery.rs:
##########
@@ -783,13 +784,16 @@ fn build_join(
 ) -> Result<Option<BuiltJoin>> {
     let mut pull_up = PullUpCorrelatedExpr::new()
         .with_in_predicate_opt(in_predicate_opt.cloned())
-        .with_exists_sub_query(in_predicate_opt.is_none());
+        .with_exists_sub_query(in_predicate_opt.is_none())
+        .with_need_handle_count_bug(true);

Review Comment:
   ```suggestion
          .with_need_handle_count_bug(in_predicate_opt.is_none());
   ```



##########
datafusion/optimizer/src/decorrelate.rs:
##########


Review Comment:
   With count-bug handling on, this arm moves a HAVING that is true on empty 
input out of the plan and into `pull_up_having_expr`. `build_join` only applies 
it on the compensated path. A `DISTINCT`, `GROUP BY` or join above the 
aggregate doesn't carry `collected_count_expr_map`, so the plan becomes a plain 
semi join and the HAVING is lost.
   
   ```sql
   -- DuckDB: 33; main: (empty, count bug); this PR: 11, 22, 44 (count(*) = 0 
is gone)
   SELECT t1_id FROM t1 WHERE EXISTS (
     SELECT DISTINCT c FROM (
       SELECT count(*) AS c FROM t2 WHERE t2.t2_id = t1.t1_id HAVING count(*) = 0
     ) s
   );
   ```
   
   When compensation can't be used, rerun the rewrite without it (main's 
behavior). `rewrite(exists)` also keeps IN on the old path:
   ```rs
   let exists = in_predicate_opt.is_none();
   let rewrite =
       |handle_count_bug: bool| -> Result<(LogicalPlan, PullUpCorrelatedExpr)> {
           let mut pull_up = PullUpCorrelatedExpr::new()
               .with_in_predicate_opt(in_predicate_opt.cloned())
               .with_exists_sub_query(exists)
               .with_need_handle_count_bug(handle_count_bug);
           let new_plan = subquery.clone().rewrite(&mut pull_up).data()?;
           Ok((new_plan, pull_up))
       };
   let (mut new_plan, mut pull_up) = rewrite(exists)?;
   let mut count_bug_compensation = pull_up
       .collected_count_expr_map
       .get(&new_plan)
       .filter(|m| !m.is_empty())
       .cloned();
   // Only the compensated join applies a HAVING moved into
   // `pull_up_having_expr`, so it needs a count map at the root.
   let compensable = pull_up.can_pull_up
       && (count_bug_compensation.is_some() || 
pull_up.pull_up_having_expr.is_none());
   if exists && !compensable {
       (new_plan, pull_up) = rewrite(false)?;
       count_bug_compensation = None;
   }
   if !pull_up.can_pull_up {
       return Ok(None);
   }
   ```



##########
datafusion/optimizer/src/decorrelate.rs:
##########
@@ -351,6 +368,49 @@ impl TreeNodeRewriter for PullUpCorrelatedExpr {
                     None
                 };
 
+                if self.pull_up_having_expr.is_some()

Review Comment:
   ```suggestion
                   if self.exists_sub_query
                       && self.pull_up_having_expr.is_some()
   ```



##########
datafusion/optimizer/src/decorrelate_predicate_subquery.rs:
##########
@@ -1063,6 +1089,116 @@ fn build_join(
     }))
 }
 
+/// Builds the join for a correlated `EXISTS` subquery whose groupless
+/// aggregate requires count-bug compensation.
+#[expect(clippy::too_many_arguments)]
+fn build_join_with_count_bug(
+    left: &LogicalPlan,
+    sub_query_alias: LogicalPlan,
+    join_filter_opt: Option<Expr>,
+    in_predicate_opt: Option<&Expr>,
+    join_type: JoinType,
+    alias: &str,
+    expr_map: &crate::decorrelate::ExprResultMap,
+    pull_up_having_expr: Option<&Expr>,
+    forces_empty_result: bool,
+) -> Result<LogicalPlan> {
+    if in_predicate_opt.is_some() {
+        return not_impl_err!(
+            "IN/NOT IN count-bug compensation for groupless aggregates is not 
implemented"
+        );
+    }
+
+    let joined = LogicalPlanBuilder::from(left.clone())
+        .join_on(sub_query_alias, JoinType::Left, join_filter_opt)?
+        .build()?;
+
+    let left_projection: Vec<Expr> = left
+        .schema()
+        .columns()
+        .into_iter()
+        .map(Expr::from)
+        .collect();
+
+    let indicator_col = Expr::Column(Column::new(Some(alias), 
UN_MATCHED_ROW_INDICATOR));
+
+    let having_arm = pull_up_having_expr.map(|f| f.clone().is_not_true());
+
+    let mut expr_rewrite = TypeCoercionRewriter::new(joined.schema());
+
+    // The subquery's HAVING clause, evaluated against an unmatched row's
+    // default aggregate values (e.g. count(*) defaults to 0, sum(x) to
+    // NULL), is usually `true`, but not once a second filter is combined
+    // into it, so it cannot be assumed.
+    let unmatched_having_default = match pull_up_having_expr {
+        Some(f) => f
+            .clone()
+            .transform_up(|e| {
+                if let Expr::Column(Column { name, .. }) = &e
+                    && let Some(default_value) = expr_map.get(name)
+                {
+                    return Ok(Transformed::yes(default_value.clone()));
+                }
+                Ok(Transformed::no(e))
+            })
+            .data()
+            .and_then(|simplified| simplified.rewrite(&mut 
expr_rewrite).data())?,
+        None => lit(true),
+    };
+
+    // EXISTS is true by default (the groupless aggregate always
+    // produces a row), unless one of three cases holds:
+    // (1) the outer row joined against an inner group whose HAVING
+    //     predicate failed,
+    // (2) an unmatched outer row's own default aggregate values fail
+    //     that same HAVING clause,
+    // (3) the subquery was truncated to zero rows unconditionally
+    //     (LIMIT 0), which is false for every outer row regardless of
+    //     whether it matched.
+    let exists_expr = if forces_empty_result {
+        lit(false)
+    } else {
+        match having_arm {
+            Some(when_expr) => {
+                when(indicator_col.clone().is_null(), unmatched_having_default)

Review Comment:
   ```suggestion
                   when(indicator_col.clone().is_null(), 
unmatched_having_default.is_true())
   ```



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