andygrove commented on code in PR #6076:
URL: https://github.com/apache/datafusion-comet/pull/6076#discussion_r4150537612
##########
spark/src/test/scala/org/apache/comet/exec/CometAggregateSuite.scala:
##########
@@ -2343,6 +2343,131 @@ class CometAggregateSuite extends CometTestBase with
AdaptiveSparkPlanHelper {
}
}
+ test("statistical aggregates with large nearby values") {
+ withSQLConf(
+ SQLConf.ADAPTIVE_EXECUTION_ENABLED.key -> "false",
+ SQLConf.SHUFFLE_PARTITIONS.key -> "1",
+ "spark.sql.files.minPartitionNum" -> "1",
+ CometConf.COMET_SHUFFLE_ENABLED.key -> "true",
+ CometConf.COMET_SHUFFLE_MODE.key -> "native") {
+ for (values <- Seq(Seq(1e16, 1e16 + 2), Seq(1e16 + 2, 1e16), Seq(-1e16,
-1e16 - 2))) {
+ // One ordered file keeps both values in the same partial accumulator.
Splitting
+ // them across files would only exercise merging two single-row states.
+ withTempPath { path =>
+ (Seq(Some(values.head), None, Some(values.last)))
+ .map(v => (0, v))
+ .toDF("g", "v")
+ .coalesce(1)
+ .write
+ .parquet(path.getCanonicalPath)
+ withParquetTable(path.getCanonicalPath, "large_moments") {
+ for (groupBy <- Seq("", " GROUP BY g")) {
+ val query = "SELECT var_pop(v), var_samp(v), stddev_pop(v),
stddev_samp(v) " +
+ "FROM large_moments" + groupBy
+ val (_, cometPlan) = checkSparkAnswerAndOperator(query)
+ val aggregates = cometPlan.collect { case a:
CometHashAggregateExec => a }
+ assert(aggregates.exists(_.modes.contains(Partial)))
+ assert(aggregates.exists(_.modes.contains(Final)))
+ checkAnswer(sql(query), Seq(Row(1.0, 2.0, 1.0, math.sqrt(2.0))))
+
+ // CORR and REGR_R2 use PearsonCorrelation's update, while
REGR_SXX/SYY
+ // and the variance used by slope/intercept follow
CentralMomentAgg.
+ checkSparkAnswerWithTolAndNumOfAggregates(
+ "SELECT corr(v, v), regr_r2(v, v), regr_sxx(v, v), regr_syy(v,
v), " +
+ "regr_slope(v, v), regr_intercept(v, v) FROM large_moments"
+ groupBy,
+ 2)
+ }
+ }
+ }
+ }
+ }
+ }
+
+ test("statistical aggregates merge large nearby values across partitions") {
+ withSQLConf(
+ SQLConf.ADAPTIVE_EXECUTION_ENABLED.key -> "false",
+ SQLConf.SHUFFLE_PARTITIONS.key -> "1",
+ SQLConf.FILES_MAX_PARTITION_BYTES.key -> "1048576",
+ SQLConf.FILES_OPEN_COST_IN_BYTES.key -> "1048576",
+ CometConf.COMET_SHUFFLE_ENABLED.key -> "true",
+ CometConf.COMET_SHUFFLE_MODE.key -> "native") {
+ withTempPath { path =>
+ // Two constant-valued files produce separate partials with zero M2.
The
+ // old merge returns 576 instead of 1024, regardless of which partial
arrives first.
+ for (value <- Seq(1e17 - 96, 1e17 - 32)) {
+ (Seq.fill(3)((0, Option(value))) ++ Seq((0, None), (1, None)))
+ .toDF("g", "v")
+ .coalesce(1)
+ .write
+ .mode("append")
+ .parquet(path.getCanonicalPath)
+ }
+ withParquetTable(path.getCanonicalPath, "merged_moments") {
+ assert(spark.table("merged_moments").rdd.getNumPartitions == 2)
+ for (groupBy <- Seq("", " GROUP BY g")) {
+ val query = "SELECT var_pop(v), var_samp(v), stddev_pop(v),
stddev_samp(v) " +
+ "FROM merged_moments" + groupBy
+ val (_, cometPlan) = checkSparkAnswerAndOperator(query)
+ val aggregates = cometPlan.collect { case a:
CometHashAggregateExec => a }
+ assert(aggregates.exists(_.modes.contains(Partial)))
+ assert(aggregates.exists(_.modes.contains(Final)))
+ val expected = Seq(Row(1024.0, 1228.8, 32.0, math.sqrt(1228.8))) ++
+ (if (groupBy.isEmpty) Seq.empty else Seq(Row(null, null, null,
null)))
+ checkAnswer(sql(query), expected)
+ checkSparkAnswerWithTolAndNumOfAggregates(
+ "SELECT covar_pop(v, -v), covar_samp(v, -v), corr(v, -v),
regr_r2(v, -v), " +
+ "regr_sxx(v, -v), regr_syy(v, -v), regr_sxy(v, -v), " +
+ "regr_slope(v, -v), regr_intercept(v, -v) FROM merged_moments"
+ groupBy,
+ 2)
+ }
+ }
+ }
+ }
+ }
+
+ test("statistical aggregates merge fractional constants across partitions") {
Review Comment:
The merge fix also resolves #6481's seven aggregate results with ANSI off:
exactly `NULL, 0.0, 0.0, 0.0, 0.0, 0.0, 0.0`. Validation passed for grouped and
ungrouped queries and both native merge orders. Could you extend this test to
cover those aggregates exactly and reference #6481?
ANSI on still returns NULL for `corr` where Spark raises `DIVIDE_BY_ZERO`,
so that remaining behavior needs a linked follow-up before closing the issue.
The fork CI note can also be updated: its run at d89c385 finished green.
##########
spark/src/test/scala/org/apache/comet/exec/CometAggregateSuite.scala:
##########
@@ -2343,6 +2343,131 @@ class CometAggregateSuite extends CometTestBase with
AdaptiveSparkPlanHelper {
}
}
+ test("statistical aggregates with large nearby values") {
+ withSQLConf(
+ SQLConf.ADAPTIVE_EXECUTION_ENABLED.key -> "false",
+ SQLConf.SHUFFLE_PARTITIONS.key -> "1",
+ "spark.sql.files.minPartitionNum" -> "1",
+ CometConf.COMET_SHUFFLE_ENABLED.key -> "true",
+ CometConf.COMET_SHUFFLE_MODE.key -> "native") {
+ for (values <- Seq(Seq(1e16, 1e16 + 2), Seq(1e16 + 2, 1e16), Seq(-1e16,
-1e16 - 2))) {
+ // One ordered file keeps both values in the same partial accumulator.
Splitting
+ // them across files would only exercise merging two single-row states.
+ withTempPath { path =>
+ (Seq(Some(values.head), None, Some(values.last)))
+ .map(v => (0, v))
+ .toDF("g", "v")
+ .coalesce(1)
+ .write
+ .parquet(path.getCanonicalPath)
+ withParquetTable(path.getCanonicalPath, "large_moments") {
+ for (groupBy <- Seq("", " GROUP BY g")) {
+ val query = "SELECT var_pop(v), var_samp(v), stddev_pop(v),
stddev_samp(v) " +
+ "FROM large_moments" + groupBy
+ val (_, cometPlan) = checkSparkAnswerAndOperator(query)
+ val aggregates = cometPlan.collect { case a:
CometHashAggregateExec => a }
+ assert(aggregates.exists(_.modes.contains(Partial)))
+ assert(aggregates.exists(_.modes.contains(Final)))
+ checkAnswer(sql(query), Seq(Row(1.0, 2.0, 1.0, math.sqrt(2.0))))
+
+ // CORR and REGR_R2 use PearsonCorrelation's update, while
REGR_SXX/SYY
+ // and the variance used by slope/intercept follow
CentralMomentAgg.
+ checkSparkAnswerWithTolAndNumOfAggregates(
Review Comment:
The scalar `corr` evaluation difference is confirmed. This nearby-value test
returns `0.9999999999999998` versus Spark's `1.0`, hidden by the `1e-6`
tolerance. On `[1e100, 2e100]`, scalar Comet returns `1.0000000000000002`
versus Spark/grouped Comet's `0.0`.
`CorrelationAccumulator::evaluate` divides normalized covariance by the two
standard deviations, while Spark and the grouped path use `ck / sqrt(m2_1 *
m2_2)`. This predates the PR. Could you either align scalar evaluation and make
`corr`'s assertion exact, or link a follow-up explicitly separating this from
the update/merge fix?
--
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]