[
https://issues.apache.org/jira/browse/SPARK-58888?page=com.atlassian.jira.plugin.system.issuetabpanels:all-tabpanel
]
Xiaoxuan Li updated SPARK-58888:
--------------------------------
Labels: Correctness (was: )
> OptimizeExpand turns an empty input into one row, making COUNT(DISTINCT)
> return 1 instead of 0
> ----------------------------------------------------------------------------------------------
>
> Key: SPARK-58888
> URL: https://issues.apache.org/jira/browse/SPARK-58888
> Project: Spark
> Issue Type: Bug
> Components: SQL
> Affects Versions: 4.2.0
> Reporter: Xiaoxuan Li
> Priority: Major
> Labels: Correctness
>
> This is 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.
> Distinct of a constant over an empty input returns 1 instead of 0. No
> constant literal has to appear in the query text: it is enough that the
> optimizer can prove the distinct columns constant, which it does for the
> common pattern of tagging constant dimensions in a subquery, view, or UNION
> ALL branch.
> ----
> *Repro*
> {code:sql}
> -- an empty source, e.g. a daily partition with no rows yet
> CREATE OR REPLACE TEMP VIEW events AS
> SELECT * FROM VALUES (1, 'x', 'y') AS t(id, c1, c2) WHERE 1 = 0;
> SET spark.sql.optimizer.optimizeExpandRatio=2;
> SELECT COUNT(DISTINCT region), COUNT(DISTINCT channel)
> FROM (SELECT 'US' AS region, 'web' AS channel, id FROM events);
> -- returns [1,1], expected [0,0]
> {code}
> Minimal form of the same thing:
> {code:sql}
> SELECT COUNT(DISTINCT 1), COUNT(DISTINCT 2) FROM events;
> -- returns [1,1], expected [0,0]
> {code}
> Excluding the rule by name gives the correct answer, so the rule is the cause
> and not the surrounding rewrite:
> {code:sql}
> SET
> spark.sql.optimizer.excludedRules=org.apache.spark.sql.catalyst.optimizer.OptimizeExpand;
> -- returns [0,0]
> {code}
> ----
> *Cause*
> {\{OptimizeExpand}} collects the pre-aggregate's grouping expressions from
> the attributes the \{{Expand}} references (\{{collectPreAggGroupBy}}):
> {code:java}
> expand.child.output.filter(expand.references.contains)
> {code}
> With no \{{GROUP BY}} and every distinct aggregate over a constant,
> \{{Expand.references}} is empty, so the inserted node is \{{Aggregate(Nil,
> Nil, child)}} -- a *global* aggregate, which emits one row even for empty
> input. The \{{Expand}} then sees one row where it previously saw none:
> {code:java}
> Aggregate [count(1) FILTER (gid=2) AS count(DISTINCT region), count(1) FILTER
> (gid=1) AS count(DISTINCT channel)]
> +- Aggregate ['web', 'US', gid], [gid]
> +- Expand [[web, null, 1], [null, US, 2]], ['web', 'US', gid]
> +- Aggregate <- inserted by OptimizeExpand: no
> grouping, no aggregates
> +- LocalRelation <empty>
> {code}
> Note this is the same empty-table failure that
> \{{RewriteDistinctAggregates.mustRewrite}} exists to prevent. Its comment
> names \{{SELECT COUNT(DISTINCT 1)}} and "wrong results when working with
> empty tables" as the reason the rewrite is forced for this shape;
> \{{OptimizeExpand}} then undoes that protection.
> ----
> *Trigger conditions*
> The distinct arguments only need to be _foldable_, not literal, so ordinary
> column references reach this state once constant projections are folded
> through:
> ||shape||result||
> |\{{COUNT(DISTINCT c1), COUNT(DISTINCT c2)}} over constant-projected
> columns|*wrong*|
> |\{{COUNT(DISTINCT 1), COUNT(DISTINCT 2)}}|*wrong*|
> |\{{COUNT(DISTINCT 1+1), COUNT(DISTINCT CAST('2020-01-01' AS DATE))}}|*wrong*|
> |three or more constant distincts|*wrong*|
> |non-empty input|correct (the answer is 1 either way)|
> |one distinct argument is a real column|correct (the column enters the
> pre-aggregate's GROUP BY)|
> |query has a \{{GROUP BY}}|correct|
> So all of the following must hold: no \{{GROUP BY}}, at least two distinct
> groups (fewer than two leaves the \{{Expand}} below the minimum threshold of
> 2), every distinct argument foldable, empty input, and the conf explicitly
> set.
> ----
> *Relationship to SPARK-58387*
> Independent. There the inner \{{Aggregate}} computes a duplicate-sensitive
> aggregate; here it computes no aggregate at all, so the duplicate-sensitivity
> guard proposed in SPARK-58387 does not reject this plan. Reproduced on a
> build that already carries that fix.
> *Suggested fix*
> Reject the rewrite when the pre-aggregate's grouping would be empty, since a
> global aggregate is not a row-preserving de-duplication -- e.g. \{{if
> (preAggGroupBy.isEmpty) return false}} in \{{canOptimize}}.
--
This message was sent by Atlassian Jira
(v8.20.10#820010)
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]