[
https://issues.apache.org/jira/browse/SPARK-58387?page=com.atlassian.jira.plugin.system.issuetabpanels:all-tabpanel
]
Josh Rosen updated SPARK-58387:
-------------------------------
Description:
This is a report of a correctness bug in the {{OptimizeExpand}} rule introduced
by SPARK-56315 (released in 4.2.0). The rule is gated by the internal conf
{{spark.sql.optimizer.optimizeExpandRatio}} (default {{-1}} = disabled), so
hitting it requires explicit opt-in.
----
Background:
* queries with multiple {{COUNT(DISTINCT)}} on different columns are rewritten
by {{RewriteDistinctAggregates}} into an {{Expand}} that duplicates each input
row once per distinct group.
* {{OptimizeExpand}} reduces that amplification by inserting a de-duplicating
{{Aggregate}} on (grouping keys + distinct columns) beneath the {{{}Expand{}}}.
* The de-duplication is sound only for pure distinct aggregates (collapsing
duplicate rows changes the value of any duplicate-sensitive aggregate computed
alongside) and the [conf's own
doc|https://github.com/apache/spark/blob/5ce51a624421acceec378e0170da24a8fc2a96eb/sql/catalyst/src/main/scala/org/apache/spark/sql/internal/SQLConf.scala#L3003-L3020]
states exactly that precondition: _"Only applies to pure distinct aggregates
without non-distinct aggregates or FILTER clauses."_
The bug: the guard implementing that precondition tests {_}attributes{_}, not
_aggregate functions_
([OptimizeExpand.scala|https://github.com/apache/spark/blob/5ce51a624421acceec378e0170da24a8fc2a96eb/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/optimizer/OptimizeExpand.scala#L106]):
{code:java}
val innerGroupByAttrs =
AttributeSet(innerAgg.groupingExpressions.flatMap(_.references))
if (!expand.producedAttributes.subsetOf(innerGroupByAttrs)) return false {code}
The reasoning is that a non-distinct aggregate would force an Expand output
column outside the inner GROUP BY. That only detects aggregates to which
{{RewriteDistinctAggregates}} assigns a dedicated Expand slot.
Any aggregate that reaches the inner Aggregate as {{count(1)}} —
{{{}COUNT(1){}}}, {{{}COUNT(*){}}}, and {{COUNT(col)}} of a non-nullable column
(which the rewrite normalizes to {{{}count(1){}}}) — references no attribute at
all and always passes the guard:
||SQL shape||after RewriteDistinctAggregates||dedicated Expand
slot?||guard||result||
|{{COUNT(1)}} / {{COUNT(*)}}|{{count(1)}}|no|passes|*wrong*|
|{{{}COUNT(a){}}}, {{a}} non-nullable|normalized to
{{count(1)}}|no|passes|*wrong*|
|{{{}COUNT(a){}}}, {{a}} nullable|{{count(a)}} over its own
slot|yes|rejects|correct|
|{{SUM(a)}}|{{sum(a)}} over its own slot|yes|rejects|correct|
{*}Bug repro{*}:
{code:java}
CREATE TABLE oe USING parquet AS SELECT * FROM VALUES
(1,5,7),(1,5,7),(1,5,7),(1,6,8),(2,9,9) AS t(k,a,b);
SET spark.sql.optimizer.optimizeExpandRatio=2;
SELECT k, COUNT(DISTINCT a), COUNT(DISTINCT b), COUNT(1) FROM oe GROUP BY k
ORDER BY k;
-- outputs: [1,2,2,2], [2,1,1,1] WRONG (expected [1,2,2,4]: k=1 has four
rows) {code}
{{k=1}} has four rows (three copies of {{(1,5,7)}} plus {{{}(1,6,8){}}}), so
{{COUNT(1)}} must be 4; the returned 2 is the number of distinct {{(k,a,b)}}
rows.
{{COUNT(*)}} behaves identically. With the conf unset (or the rule excluded)
the answer is correct.
The optimized plan shows the inserted pre-aggregate (last {{{}Aggregate [k, a,
b]{}}}) feeding the Expand whose downstream {{count(1)}} was supposed to count
base rows:
{code:java}
Aggregate [k], [k, count(a) FILTER (gid=1), count(b) FILTER (gid=2),
coalesce(first(count(1)) FILTER (gid=0), 0)]
+- Aggregate [k, a, b, gid], [k, a, b, gid, count(1)]
+- Expand [[k,null,null,0], [k,a,null,1], [k,null,b,2]], [k, a, b, gid]
+- Aggregate [k, a, b], [k, a, b] <- inserted by OptimizeExpand
+- Relation oe[k,a,b] parquet {code}
With the rule excluded, the Expand's child is the relation itself and
{{count(1)}} is 4.
To fix this: check the inner aggregate's aggregate expressions rather than its
grouping references — reject any {{AggregateExpression}} that is non-distinct
and not duplicate-agnostic (cf. {{{}EliminateDistinct.isDuplicateAgnostic{}}},
the check {{RemoveRedundantAggregates}} uses for the same soundness question).
was:
This is a report of a correctness bug in the {{OptimizeExpand}} rule introduced
by SPARK-56315 (released in 4.2.0). The rule is gated by the internal conf
{{spark.sql.optimizer.optimizeExpandRatio}} (default {{-1}} = disabled), so
hitting it requires explicit opt-in.
----
Background:
* queries with multiple {{COUNT(DISTINCT)}} on different columns are rewritten
by {{RewriteDistinctAggregates}} into an {{Expand}} that duplicates each input
row once per distinct group.
* {{OptimizeExpand}} reduces that amplification by inserting a de-duplicating
{{Aggregate}} on (grouping keys + distinct columns) beneath the {{{}Expand{}}}.
* The de-duplication is sound only for pure distinct aggregates (collapsing
duplicate rows changes the value of any duplicate-sensitive aggregate computed
alongside) and the [conf's own
doc|https://github.com/apache/spark/blob/5ce51a624421acceec378e0170da24a8fc2a96eb/sql/catalyst/src/main/scala/org/apache/spark/sql/internal/SQLConf.scala#L3003-L3020]
states exactly that precondition: _"Only applies to pure distinct aggregates
without non-distinct aggregates or FILTER clauses."_
The bug: the guard implementing that precondition tests {_}attributes{_}, not
_aggregate functions_
([OptimizeExpand.scala|https://github.com/apache/spark/blob/5ce51a624421acceec378e0170da24a8fc2a96eb/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/optimizer/OptimizeExpand.scala#L106]):
{code:java}
val innerGroupByAttrs =
AttributeSet(innerAgg.groupingExpressions.flatMap(_.references))
if (!expand.producedAttributes.subsetOf(innerGroupByAttrs)) return false {code}
The reasoning is that a non-distinct aggregate would force an Expand output
column outside the inner GROUP BY. That only detects aggregates to which
{{RewriteDistinctAggregates}} assigns a dedicated Expand slot.
Any aggregate that reaches the inner Aggregate as {{count(1)}} —
{{{}COUNT(1){}}}, {{{}COUNT(*){}}}, and {{COUNT(col)}} of a non-nullable column
(which the rewrite normalizes to {{{}count(1){}}}) — references no attribute at
all and always passes the guard:
||SQL shape||after RewriteDistinctAggregates||dedicated Expand
slot?||guard||result||
|{{COUNT(1)}} / {{COUNT(*)}}|{{count(1)}}|no|passes|*wrong*|
|{{{}COUNT(a){}}}, {{a}} non-nullable|normalized to
{{count(1)}}|no|passes|*wrong*|
|{{{}COUNT(a){}}}, {{a}} nullable|{{count(a)}} over its own
slot|yes|rejects|correct|
|{{SUM(a)}}|{{sum(a)}} over its own slot|yes|rejects|correct|
{*}Bug repro{*}:
{code:java}
CREATE TABLE oe USING parquet AS SELECT * FROM VALUES
(1,5,7),(1,5,7),(1,5,7),(1,6,8),(2,9,9) AS t(k,a,b);
SET spark.sql.optimizer.optimizeExpandRatio=2;
SELECT k, COUNT(DISTINCT a), COUNT(DISTINCT b), COUNT(1) FROM oe GROUP BY k
ORDER BY k;
-- outputs: [1,2,2,2], [2,1,1,1] WRONG (expected [1,2,2,4]: k=1 has four
rows) {code}
{{}}
{{}}
{{k=1}} has four rows (three copies of {{(1,5,7)}} plus {{{}(1,6,8){}}}), so
{{COUNT(1)}} must be 4; the returned 2 is the number of distinct {{(k,a,b)}}
rows.
{{COUNT(*)}} behaves identically. With the conf unset (or the rule excluded)
the answer is correct.
The optimized plan shows the inserted pre-aggregate (last {{{}Aggregate [k, a,
b]{}}}) feeding the Expand whose downstream {{count(1)}} was supposed to count
base rows:
{code:java}
Aggregate [k], [k, count(a) FILTER (gid=1), count(b) FILTER (gid=2),
coalesce(first(count(1)) FILTER (gid=0), 0)]
+- Aggregate [k, a, b, gid], [k, a, b, gid, count(1)]
+- Expand [[k,null,null,0], [k,a,null,1], [k,null,b,2]], [k, a, b, gid]
+- Aggregate [k, a, b], [k, a, b] <- inserted by OptimizeExpand
+- Relation oe[k,a,b] parquet {code}
{{}}
{{}}
With the rule excluded, the Expand's child is the relation itself and
{{count(1)}} is 4.
To fix this: check the inner aggregate's aggregate expressions rather than its
grouping references — reject any {{AggregateExpression}} that is non-distinct
and not duplicate-agnostic (cf. {{{}EliminateDistinct.isDuplicateAgnostic{}}},
the check {{RemoveRedundantAggregates}} uses for the same soundness question).
> OptimizeExpand inserts a de-duplicating pre-aggregate under a non-distinct
> COUNT, returning the distinct-row count
> ------------------------------------------------------------------------------------------------------------------
>
> Key: SPARK-58387
> URL: https://issues.apache.org/jira/browse/SPARK-58387
> Project: Spark
> Issue Type: Bug
> Components: SQL
> Affects Versions: 4.2.0
> Reporter: Josh Rosen
> Priority: Major
> Labels: correctness
>
> This is a report of a correctness bug in the {{OptimizeExpand}} rule
> introduced by SPARK-56315 (released in 4.2.0). The rule is gated by the
> internal conf {{spark.sql.optimizer.optimizeExpandRatio}} (default {{-1}} =
> disabled), so hitting it requires explicit opt-in.
>
> ----
>
> Background:
> * queries with multiple {{COUNT(DISTINCT)}} on different columns are
> rewritten by {{RewriteDistinctAggregates}} into an {{Expand}} that duplicates
> each input row once per distinct group.
> * {{OptimizeExpand}} reduces that amplification by inserting a
> de-duplicating {{Aggregate}} on (grouping keys + distinct columns) beneath
> the {{{}Expand{}}}.
> * The de-duplication is sound only for pure distinct aggregates (collapsing
> duplicate rows changes the value of any duplicate-sensitive aggregate
> computed alongside) and the [conf's own
> doc|https://github.com/apache/spark/blob/5ce51a624421acceec378e0170da24a8fc2a96eb/sql/catalyst/src/main/scala/org/apache/spark/sql/internal/SQLConf.scala#L3003-L3020]
> states exactly that precondition: _"Only applies to pure distinct aggregates
> without non-distinct aggregates or FILTER clauses."_
>
> The bug: the guard implementing that precondition tests {_}attributes{_}, not
> _aggregate functions_
> ([OptimizeExpand.scala|https://github.com/apache/spark/blob/5ce51a624421acceec378e0170da24a8fc2a96eb/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/optimizer/OptimizeExpand.scala#L106]):
> {code:java}
> val innerGroupByAttrs =
> AttributeSet(innerAgg.groupingExpressions.flatMap(_.references))
> if (!expand.producedAttributes.subsetOf(innerGroupByAttrs)) return false
> {code}
>
> The reasoning is that a non-distinct aggregate would force an Expand output
> column outside the inner GROUP BY. That only detects aggregates to which
> {{RewriteDistinctAggregates}} assigns a dedicated Expand slot.
> Any aggregate that reaches the inner Aggregate as {{count(1)}} —
> {{{}COUNT(1){}}}, {{{}COUNT(*){}}}, and {{COUNT(col)}} of a non-nullable
> column (which the rewrite normalizes to {{{}count(1){}}}) — references no
> attribute at all and always passes the guard:
> ||SQL shape||after RewriteDistinctAggregates||dedicated Expand
> slot?||guard||result||
> |{{COUNT(1)}} / {{COUNT(*)}}|{{count(1)}}|no|passes|*wrong*|
> |{{{}COUNT(a){}}}, {{a}} non-nullable|normalized to
> {{count(1)}}|no|passes|*wrong*|
> |{{{}COUNT(a){}}}, {{a}} nullable|{{count(a)}} over its own
> slot|yes|rejects|correct|
> |{{SUM(a)}}|{{sum(a)}} over its own slot|yes|rejects|correct|
>
> {*}Bug repro{*}:
> {code:java}
> CREATE TABLE oe USING parquet AS SELECT * FROM VALUES
> (1,5,7),(1,5,7),(1,5,7),(1,6,8),(2,9,9) AS t(k,a,b);
> SET spark.sql.optimizer.optimizeExpandRatio=2;
> SELECT k, COUNT(DISTINCT a), COUNT(DISTINCT b), COUNT(1) FROM oe GROUP BY k
> ORDER BY k;
> -- outputs: [1,2,2,2], [2,1,1,1] WRONG (expected [1,2,2,4]: k=1 has four
> rows) {code}
> {{k=1}} has four rows (three copies of {{(1,5,7)}} plus {{{}(1,6,8){}}}), so
> {{COUNT(1)}} must be 4; the returned 2 is the number of distinct {{(k,a,b)}}
> rows.
> {{COUNT(*)}} behaves identically. With the conf unset (or the rule excluded)
> the answer is correct.
>
> The optimized plan shows the inserted pre-aggregate (last {{{}Aggregate [k,
> a, b]{}}}) feeding the Expand whose downstream {{count(1)}} was supposed to
> count base rows:
> {code:java}
> Aggregate [k], [k, count(a) FILTER (gid=1), count(b) FILTER (gid=2),
> coalesce(first(count(1)) FILTER (gid=0), 0)]
> +- Aggregate [k, a, b, gid], [k, a, b, gid, count(1)]
> +- Expand [[k,null,null,0], [k,a,null,1], [k,null,b,2]], [k, a, b, gid]
> +- Aggregate [k, a, b], [k, a, b] <- inserted by OptimizeExpand
> +- Relation oe[k,a,b] parquet {code}
> With the rule excluded, the Expand's child is the relation itself and
> {{count(1)}} is 4.
>
> To fix this: check the inner aggregate's aggregate expressions rather than
> its grouping references — reject any {{AggregateExpression}} that is
> non-distinct and not duplicate-agnostic (cf.
> {{{}EliminateDistinct.isDuplicateAgnostic{}}}, the check
> {{RemoveRedundantAggregates}} uses for the same soundness question).
--
This message was sent by Atlassian Jira
(v8.20.10#820010)
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]