[
https://issues.apache.org/jira/browse/SPARK-58385?page=com.atlassian.jira.plugin.system.issuetabpanels:all-tabpanel
]
Josh Rosen updated SPARK-58385:
-------------------------------
Description:
This is a report of a correctness bug present in Spark 2.0.0 and all later
versions.
[{{LimitPushDown.maybePushLocalLimit}}|https://github.com/apache/spark/blob/5ce51a624421acceec378e0170da24a8fc2a96eb/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/optimizer/Optimizer.scala#L916]
decides whether a pushed-down limit may replace a child's existing limit by
inspecting {{{}plan.maxRowsPerPartition{}}}:
{code:java}
private def maybePushLocalLimit(limitExp: Expression, plan: LogicalPlan):
LogicalPlan = {
(limitExp, plan.maxRowsPerPartition) match {
case (IntegerLiteral(newLimit), Some(childMaxRows)) if newLimit <
childMaxRows =>
// If the child has a cap on max rows per partition and the cap is
larger than
// the new limit, put a new LocalLimit there.
LocalLimit(limitExp, stripGlobalLimitIfPresent(plan))
case (_, None) =>
// If the child has no cap, put the new LocalLimit.
LocalLimit(limitExp, stripGlobalLimitIfPresent(plan))
case _ =>
// Otherwise, don't put a new LocalLimit.
plan
}
}{code}
* The first arm strips a child {{GlobalLimit}} only under a guard
({{{}newLimit < childMaxRows{}}}).
* The {{(_, None)}} arm strips unconditionally, reading {{None}} as "the child
has no cap" — but {{None}} also means "the cap is statically unknown".
*
[{{GlobalLimit.maxRows}}|https://github.com/apache/spark/blob/5ce51a624421acceec378e0170da24a8fc2a96eb/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/plans/logical/basicLogicalOperators.scala#L1896-L1901]
/ {{maxRowsPerPartition}} recognize only {{{}IntegerLiteral{}}}, and a limit
written as a foldable expression ({{{}LIMIT 2+3{}}}) is still {{Add(2, 3)}}
when {{LimitPushDown}} runs: it is
[listed|https://github.com/apache/spark/blob/5ce51a624421acceec378e0170da24a8fc2a96eb/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/optimizer/Optimizer.scala#L112]
before {{ConstantFolding}} in the same operator-optimization rule set, so on
the first fixed-point pass the expression is unfolded, {{maxRowsPerPartition}}
returns {{{}None{}}}, and the rule deletes the semantic {{{}GlobalLimit{}}},
replacing a global row cap with a per-partition one.
* Foldable non-literal LIMIT expressions are valid SQL:
{{checkLimitLikeClause}} accepts any foldable integer expression; they arise
naturally from templated/generated SQL that renders arithmetic into the query
text, e.g. a pagination template producing {{{}LIMIT 20*5{}}}.
{*}Repro{*}:
{code:java}
CREATE TABLE la USING parquet AS SELECT 1 AS id;
CREATE TABLE lb USING parquet AS SELECT id FROM range(0, 40, 1, 4);
SELECT count(*) FROM (SELECT * FROM la CROSS JOIN (SELECT * FROM lb LIMIT 2+3)
LIMIT 100);
-- outputs: 10 WRONG (expected: 5){code}
This returns an incorrect answer because the larger outer limit strips the
nested global limit during its pushdown. If we manually fold the inner limit to
{{5}} then the query outputs the expected 5 rows.
{*}Fix{*}:
Treat {{None}} as "unknown", not "uncapped": in the {{(_, None)}} arm, never
strip a {{GlobalLimit}} child (push the {{LocalLimit}} beneath it, or leave the
plan unchanged). Alternatively (or additionally), make {{GlobalLimit.maxRows}}
evaluate foldable limit expressions so the guarded arm applies. Stripping is
only sound when the child's cap is known and the guard {{newLimit <
childMaxRows}} holds, as the literal arm already implements.
was:
This is a report of a correctness bug present in Spark 2.0.0 and all later
versions.
[{{LimitPushDown.maybePushLocalLimit}}|https://github.com/apache/spark/blob/5ce51a624421acceec378e0170da24a8fc2a96eb/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/optimizer/Optimizer.scala#L916]
decides whether a pushed-down limit may replace a child's existing limit by
inspecting {{{}plan.maxRowsPerPartition{}}}:
{code:java}
private def maybePushLocalLimit(limitExp: Expression, plan: LogicalPlan):
LogicalPlan = {
(limitExp, plan.maxRowsPerPartition) match {
case (IntegerLiteral(newLimit), Some(childMaxRows)) if newLimit <
childMaxRows =>
// If the child has a cap on max rows per partition and the cap is
larger than
// the new limit, put a new LocalLimit there.
LocalLimit(limitExp, stripGlobalLimitIfPresent(plan))
case (_, None) =>
// If the child has no cap, put the new LocalLimit.
LocalLimit(limitExp, stripGlobalLimitIfPresent(plan))
case _ =>
// Otherwise, don't put a new LocalLimit.
plan
}
}{code}
* The first arm strips a child {{GlobalLimit}} only under a guard
({{{}newLimit < childMaxRows{}}}).
* The {{(_, None)}} arm strips unconditionally, reading {{None}} as "the child
has no cap" — but {{None}} also means "the cap is statically unknown".
*
[{{GlobalLimit.maxRows}}|https://github.com/apache/spark/blob/5ce51a624421acceec378e0170da24a8fc2a96eb/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/plans/logical/basicLogicalOperators.scala#L1896-L1901]
/ {{maxRowsPerPartition}} recognize only {{{}IntegerLiteral{}}}, and a limit
written as a foldable expression ({{{}LIMIT 2+3{}}}) is still {{Add(2, 3)}}
when {{LimitPushDown}} runs: it is
[listed|https://github.com/apache/spark/blob/5ce51a624421acceec378e0170da24a8fc2a96eb/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/optimizer/Optimizer.scala#L112]
before {{ConstantFolding}} in the same operator-optimization rule set, so on
the first fixed-point pass the expression is unfolded, {{maxRowsPerPartition}}
returns {{{}None{}}}, and the rule deletes the semantic {{{}GlobalLimit{}}},
replacing a global row cap with a per-partition one.
* Foldable non-literal LIMIT expressions are valid SQL:
{{checkLimitLikeClause}} accepts any foldable integer expression; they arise
naturally from templated/generated SQL that renders arithmetic into the query
text, e.g. a pagination template producing {{{}LIMIT 20*5{}}}.
{*}Repro{*}:
{code:java}
CREATE TABLE la USING parquet AS SELECT 1 AS id;
CREATE TABLE lb USING parquet AS SELECT id FROM range(0, 40, 1, 4);
SELECT count(*) FROM (SELECT * FROM la CROSS JOIN (SELECT * FROM lb LIMIT 2+3)
LIMIT 100);
-- outputs: 10 WRONG (expected: 5){code}
This returns an incorrect answer because the larger outer limit strips the
nested global limit. If we manually fold the inner limit to {{5}} .
{*}Fix{*}:
Treat {{None}} as "unknown", not "uncapped": in the {{(_, None)}} arm, never
strip a {{GlobalLimit}} child (push the {{LocalLimit}} beneath it, or leave the
plan unchanged). Alternatively (or additionally), make {{GlobalLimit.maxRows}}
evaluate foldable limit expressions so the guarded arm applies. Stripping is
only sound when the child's cap is known and the guard {{newLimit <
childMaxRows}} holds, as the literal arm already implements.
> LimitPushDown incorrectly removes a GlobalLimit whose limit expression is not
> a literal
> ---------------------------------------------------------------------------------------
>
> Key: SPARK-58385
> URL: https://issues.apache.org/jira/browse/SPARK-58385
> Project: Spark
> Issue Type: Bug
> Components: SQL
> Affects Versions: 4.0.0
> Reporter: Josh Rosen
> Priority: Major
> Labels: correctness
>
> This is a report of a correctness bug present in Spark 2.0.0 and all later
> versions.
> [{{LimitPushDown.maybePushLocalLimit}}|https://github.com/apache/spark/blob/5ce51a624421acceec378e0170da24a8fc2a96eb/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/optimizer/Optimizer.scala#L916]
> decides whether a pushed-down limit may replace a child's existing limit by
> inspecting {{{}plan.maxRowsPerPartition{}}}:
> {code:java}
> private def maybePushLocalLimit(limitExp: Expression, plan: LogicalPlan):
> LogicalPlan = {
> (limitExp, plan.maxRowsPerPartition) match {
> case (IntegerLiteral(newLimit), Some(childMaxRows)) if newLimit <
> childMaxRows =>
> // If the child has a cap on max rows per partition and the cap is
> larger than
> // the new limit, put a new LocalLimit there.
> LocalLimit(limitExp, stripGlobalLimitIfPresent(plan))
> case (_, None) =>
> // If the child has no cap, put the new LocalLimit.
> LocalLimit(limitExp, stripGlobalLimitIfPresent(plan))
> case _ =>
> // Otherwise, don't put a new LocalLimit.
> plan
> }
> }{code}
> * The first arm strips a child {{GlobalLimit}} only under a guard
> ({{{}newLimit < childMaxRows{}}}).
> * The {{(_, None)}} arm strips unconditionally, reading {{None}} as "the
> child has no cap" — but {{None}} also means "the cap is statically unknown".
> *
> [{{GlobalLimit.maxRows}}|https://github.com/apache/spark/blob/5ce51a624421acceec378e0170da24a8fc2a96eb/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/plans/logical/basicLogicalOperators.scala#L1896-L1901]
> / {{maxRowsPerPartition}} recognize only {{{}IntegerLiteral{}}}, and a limit
> written as a foldable expression ({{{}LIMIT 2+3{}}}) is still {{Add(2, 3)}}
> when {{LimitPushDown}} runs: it is
> [listed|https://github.com/apache/spark/blob/5ce51a624421acceec378e0170da24a8fc2a96eb/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/optimizer/Optimizer.scala#L112]
> before {{ConstantFolding}} in the same operator-optimization rule set, so on
> the first fixed-point pass the expression is unfolded,
> {{maxRowsPerPartition}} returns {{{}None{}}}, and the rule deletes the
> semantic {{{}GlobalLimit{}}}, replacing a global row cap with a per-partition
> one.
> * Foldable non-literal LIMIT expressions are valid SQL:
> {{checkLimitLikeClause}} accepts any foldable integer expression; they arise
> naturally from templated/generated SQL that renders arithmetic into the query
> text, e.g. a pagination template producing {{{}LIMIT 20*5{}}}.
> {*}Repro{*}:
> {code:java}
> CREATE TABLE la USING parquet AS SELECT 1 AS id;
> CREATE TABLE lb USING parquet AS SELECT id FROM range(0, 40, 1, 4);
> SELECT count(*) FROM (SELECT * FROM la CROSS JOIN (SELECT * FROM lb LIMIT
> 2+3) LIMIT 100);
> -- outputs: 10 WRONG (expected: 5){code}
> This returns an incorrect answer because the larger outer limit strips the
> nested global limit during its pushdown. If we manually fold the inner limit
> to {{5}} then the query outputs the expected 5 rows.
> {*}Fix{*}:
> Treat {{None}} as "unknown", not "uncapped": in the {{(_, None)}} arm, never
> strip a {{GlobalLimit}} child (push the {{LocalLimit}} beneath it, or leave
> the plan unchanged). Alternatively (or additionally), make
> {{GlobalLimit.maxRows}} evaluate foldable limit expressions so the guarded
> arm applies. Stripping is only sound when the child's cap is known and the
> guard {{newLimit < childMaxRows}} holds, as the literal arm already
> implements.
--
This message was sent by Atlassian Jira
(v8.20.10#820010)
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]