From 09406e2c037d8edb619441b6438f085c5398e1f4 Mon Sep 17 00:00:00 2001 From: Takeshi Yamamuro Date: Fri, 15 Nov 2019 10:39:37 +0900 Subject: [PATCH 1/4] Fix --- .../spark/sql/catalyst/optimizer/Optimizer.scala | 1 - .../scala/org/apache/spark/sql/SubquerySuite.scala | 12 ++++++++++++ 2 files changed, 12 insertions(+), 1 deletion(-) diff --git a/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/optimizer/Optimizer.scala b/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/optimizer/Optimizer.scala index b78bdf082f33..ecba67190116 100644 --- a/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/optimizer/Optimizer.scala +++ b/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/optimizer/Optimizer.scala @@ -1006,7 +1006,6 @@ object EliminateSorts extends Rule[LogicalPlan] { case _: Min => true case _: Max => true case _: Count => true - case _: Average => true case _: CentralMomentAgg => true case _ => false } diff --git a/sql/core/src/test/scala/org/apache/spark/sql/SubquerySuite.scala b/sql/core/src/test/scala/org/apache/spark/sql/SubquerySuite.scala index c117ee7818c0..4fe5d4cafaae 100644 --- a/sql/core/src/test/scala/org/apache/spark/sql/SubquerySuite.scala +++ b/sql/core/src/test/scala/org/apache/spark/sql/SubquerySuite.scala @@ -1271,6 +1271,18 @@ class SubquerySuite extends QueryTest with SharedSparkSession { } } + test("Cannot remove sort for AVG from subquery plan because it is order-sensitive") { + val query = + """ + |SELECT k, AVG(v) FROM ( + | SELECT k, v + | FROM VALUES (1, 2), (2, 1) t(k, v) + | ORDER BY v) + |GROUP BY k + """.stripMargin + assert(getNumSortsInQuery(query) == 1) + } + test("SPARK-25482: Forbid pushdown to datasources of filters containing subqueries") { withTempView("t1", "t2") { sql("create temporary view t1(a int) using parquet") From a44f21890f6f52e80a80f31a0437e0ac99fe6b98 Mon Sep 17 00:00:00 2001 From: Takeshi Yamamuro Date: Fri, 15 Nov 2019 12:48:24 +0900 Subject: [PATCH 2/4] Fix --- .../sql/catalyst/optimizer/Optimizer.scala | 2 ++ .../org/apache/spark/sql/SubquerySuite.scala | 19 ++++++++++++++++--- 2 files changed, 18 insertions(+), 3 deletions(-) diff --git a/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/optimizer/Optimizer.scala b/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/optimizer/Optimizer.scala index ecba67190116..23dbd45cf44c 100644 --- a/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/optimizer/Optimizer.scala +++ b/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/optimizer/Optimizer.scala @@ -1007,6 +1007,8 @@ object EliminateSorts extends Rule[LogicalPlan] { case _: Max => true case _: Count => true case _: CentralMomentAgg => true + // Floating-piont Average aggregates are order-sensitive + case a: Average => !Seq(FloatType, DoubleType).exists(_.sameType(a.child.dataType)) case _ => false } diff --git a/sql/core/src/test/scala/org/apache/spark/sql/SubquerySuite.scala b/sql/core/src/test/scala/org/apache/spark/sql/SubquerySuite.scala index 4fe5d4cafaae..7c54dbf4887a 100644 --- a/sql/core/src/test/scala/org/apache/spark/sql/SubquerySuite.scala +++ b/sql/core/src/test/scala/org/apache/spark/sql/SubquerySuite.scala @@ -1271,8 +1271,21 @@ class SubquerySuite extends QueryTest with SharedSparkSession { } } - test("Cannot remove sort for AVG from subquery plan because it is order-sensitive") { - val query = + test("Cannot remove sort for floating-point AVG from subquery because it is order-sensitive") { + Seq("float", "double").foreach { typeName => + val query1 = + s""" + |SELECT k, AVG(v) FROM ( + | SELECT k, v + | FROM VALUES (1, $typeName(2.0)), (2, $typeName(1.0)) t(k, v) + | ORDER BY v) + |GROUP BY k + """.stripMargin + assert(getNumSortsInQuery(query1) == 1) + } + + // For integral aggreagtes, we can remove redundant sort in a plan + val query2 = """ |SELECT k, AVG(v) FROM ( | SELECT k, v @@ -1280,7 +1293,7 @@ class SubquerySuite extends QueryTest with SharedSparkSession { | ORDER BY v) |GROUP BY k """.stripMargin - assert(getNumSortsInQuery(query) == 1) + assert(getNumSortsInQuery(query2) == 0) } test("SPARK-25482: Forbid pushdown to datasources of filters containing subqueries") { From a1815ba2c613cc6c50577db3f40505302a4e4e13 Mon Sep 17 00:00:00 2001 From: Takeshi Yamamuro Date: Fri, 15 Nov 2019 21:42:36 +0900 Subject: [PATCH 3/4] Fix --- .../sql/catalyst/optimizer/Optimizer.scala | 8 ++--- .../org/apache/spark/sql/SubquerySuite.scala | 34 +++++++------------ 2 files changed, 17 insertions(+), 25 deletions(-) diff --git a/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/optimizer/Optimizer.scala b/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/optimizer/Optimizer.scala index 23dbd45cf44c..8c63cf17289b 100644 --- a/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/optimizer/Optimizer.scala +++ b/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/optimizer/Optimizer.scala @@ -1002,13 +1002,13 @@ object EliminateSorts extends Rule[LogicalPlan] { private def isOrderIrrelevantAggs(aggs: Seq[NamedExpression]): Boolean = { def isOrderIrrelevantAggFunction(func: AggregateFunction): Boolean = func match { - case _: Sum => true case _: Min => true case _: Max => true case _: Count => true - case _: CentralMomentAgg => true - // Floating-piont Average aggregates are order-sensitive - case a: Average => !Seq(FloatType, DoubleType).exists(_.sameType(a.child.dataType)) + // Arithmetic operations for floating-point values are order-sensitive + // (they are not associative). + case _: Sum | _: Average | _: CentralMomentAgg => + !Seq(FloatType, DoubleType).exists(_.sameType(func.children.head.dataType)) case _ => false } diff --git a/sql/core/src/test/scala/org/apache/spark/sql/SubquerySuite.scala b/sql/core/src/test/scala/org/apache/spark/sql/SubquerySuite.scala index 7c54dbf4887a..d7cc5c238f6e 100644 --- a/sql/core/src/test/scala/org/apache/spark/sql/SubquerySuite.scala +++ b/sql/core/src/test/scala/org/apache/spark/sql/SubquerySuite.scala @@ -1271,29 +1271,21 @@ class SubquerySuite extends QueryTest with SharedSparkSession { } } - test("Cannot remove sort for floating-point AVG from subquery because it is order-sensitive") { + test("Cannot remove sort for floating-point order-sensitive aggregates from subquery") { Seq("float", "double").foreach { typeName => - val query1 = - s""" - |SELECT k, AVG(v) FROM ( - | SELECT k, v - | FROM VALUES (1, $typeName(2.0)), (2, $typeName(1.0)) t(k, v) - | ORDER BY v) - |GROUP BY k - """.stripMargin - assert(getNumSortsInQuery(query1) == 1) + Seq("SUM", "AVG", "KURTOSIS", "SKEWNESS", "STDDEV_POP", "STDDEV_SAMP", + "VAR_POP", "VAR_SAMP").foreach { aggName => + val query1 = + s""" + |SELECT k, $aggName(v) FROM ( + | SELECT k, v + | FROM VALUES (1, $typeName(2.0)), (2, $typeName(1.0)) t(k, v) + | ORDER BY v) + |GROUP BY k + """.stripMargin + assert(getNumSortsInQuery(query1) == 1) + } } - - // For integral aggreagtes, we can remove redundant sort in a plan - val query2 = - """ - |SELECT k, AVG(v) FROM ( - | SELECT k, v - | FROM VALUES (1, 2), (2, 1) t(k, v) - | ORDER BY v) - |GROUP BY k - """.stripMargin - assert(getNumSortsInQuery(query2) == 0) } test("SPARK-25482: Forbid pushdown to datasources of filters containing subqueries") { From ebec3175b726828cf38ce79fd430c9d8957835ca Mon Sep 17 00:00:00 2001 From: Takeshi Yamamuro Date: Sat, 16 Nov 2019 07:16:59 +0900 Subject: [PATCH 4/4] Fix --- .../org/apache/spark/sql/catalyst/optimizer/Optimizer.scala | 4 +--- .../src/test/scala/org/apache/spark/sql/SubquerySuite.scala | 4 ++-- 2 files changed, 3 insertions(+), 5 deletions(-) diff --git a/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/optimizer/Optimizer.scala b/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/optimizer/Optimizer.scala index 8c63cf17289b..473f846c9313 100644 --- a/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/optimizer/Optimizer.scala +++ b/sql/catalyst/src/main/scala/org/apache/spark/sql/catalyst/optimizer/Optimizer.scala @@ -1002,9 +1002,7 @@ object EliminateSorts extends Rule[LogicalPlan] { private def isOrderIrrelevantAggs(aggs: Seq[NamedExpression]): Boolean = { def isOrderIrrelevantAggFunction(func: AggregateFunction): Boolean = func match { - case _: Min => true - case _: Max => true - case _: Count => true + case _: Min | _: Max | _: Count => true // Arithmetic operations for floating-point values are order-sensitive // (they are not associative). case _: Sum | _: Average | _: CentralMomentAgg => diff --git a/sql/core/src/test/scala/org/apache/spark/sql/SubquerySuite.scala b/sql/core/src/test/scala/org/apache/spark/sql/SubquerySuite.scala index d7cc5c238f6e..5020c1047f8d 100644 --- a/sql/core/src/test/scala/org/apache/spark/sql/SubquerySuite.scala +++ b/sql/core/src/test/scala/org/apache/spark/sql/SubquerySuite.scala @@ -1275,7 +1275,7 @@ class SubquerySuite extends QueryTest with SharedSparkSession { Seq("float", "double").foreach { typeName => Seq("SUM", "AVG", "KURTOSIS", "SKEWNESS", "STDDEV_POP", "STDDEV_SAMP", "VAR_POP", "VAR_SAMP").foreach { aggName => - val query1 = + val query = s""" |SELECT k, $aggName(v) FROM ( | SELECT k, v @@ -1283,7 +1283,7 @@ class SubquerySuite extends QueryTest with SharedSparkSession { | ORDER BY v) |GROUP BY k """.stripMargin - assert(getNumSortsInQuery(query1) == 1) + assert(getNumSortsInQuery(query) == 1) } } }