From bcf1a5f6c0abb554a45436832e444b906ee62495 Mon Sep 17 00:00:00 2001 From: Andy Grove Date: Sat, 3 Oct 2026 10:00:33 -0600 Subject: [PATCH] docs: regenerate release docs for 1.1.0-rc2 Output of ./dev/generate-release-docs.sh on branch-1.1 after the backports that merged since rc1: the skipPartial config row (#6488), the regr_* fallback reasons (#6489) and the array_distinct/array_union fallback reasons (#6561) on every Spark version's compatibility pages, and the regr_sxx/regr_syy Implementation cells back to the generator's placeholder. --- .../expressions/spark-3.4/aggregate.md | 30 +++++++++++++++++++ .../expressions/spark-3.4/array.md | 12 ++++++++ .../expressions/spark-3.5/aggregate.md | 30 +++++++++++++++++++ .../expressions/spark-3.5/array.md | 12 ++++++++ .../expressions/spark-4.0/aggregate.md | 30 +++++++++++++++++++ .../expressions/spark-4.0/array.md | 12 ++++++++ .../expressions/spark-4.1/aggregate.md | 30 +++++++++++++++++++ .../expressions/spark-4.1/array.md | 12 ++++++++ docs/source/user-guide/latest/configs.md | 1 + docs/source/user-guide/latest/expressions.md | 4 +-- 10 files changed, 171 insertions(+), 2 deletions(-) diff --git a/docs/source/user-guide/latest/compatibility/expressions/spark-3.4/aggregate.md b/docs/source/user-guide/latest/compatibility/expressions/spark-3.4/aggregate.md index daf7a9c00a5..22300066ce6 100644 --- a/docs/source/user-guide/latest/compatibility/expressions/spark-3.4/aggregate.md +++ b/docs/source/user-guide/latest/compatibility/expressions/spark-3.4/aggregate.md @@ -97,4 +97,34 @@ The following cases are not supported by Comet and always fall back to Spark, re - Descending order in `WITHIN GROUP (ORDER BY ... DESC)` is not supported. - Only numeric input types are supported. +## RegrIntercept + +The following incompatibilities cause `RegrIntercept` to fall back to Spark by default. Set `spark.comet.expression.RegrIntercept.allowIncompatible=true` to enable Comet acceleration despite these differences. + +- Comet merges the partial aggregates of `regr_intercept` in a different floating-point operation order from Spark. When a group's rows come from more than one partial aggregate and a variable is constant at a value that binary floating point cannot represent exactly, such as 0.1, Comet returns a wrong value where Spark returns NULL, 0.0 or 1.0 (https://github.com/apache/datafusion-comet/issues/6423) + +## RegrR2 + +The following incompatibilities cause `RegrR2` to fall back to Spark by default. Set `spark.comet.expression.RegrR2.allowIncompatible=true` to enable Comet acceleration despite these differences. + +- Comet merges the partial aggregates of `regr_r2` in a different floating-point operation order from Spark. When a group's rows come from more than one partial aggregate and a variable is constant at a value that binary floating point cannot represent exactly, such as 0.1, Comet returns a wrong value where Spark returns NULL, 0.0 or 1.0 (https://github.com/apache/datafusion-comet/issues/6423) + +## RegrReplacement + +The following incompatibilities cause `RegrReplacement` to fall back to Spark by default. Set `spark.comet.expression.RegrReplacement.allowIncompatible=true` to enable Comet acceleration despite these differences. + +- Comet merges the partial aggregates of `regr_sxx` and `regr_syy` in a different floating-point operation order from Spark. When a group's rows come from more than one partial aggregate and a variable is constant at a value that binary floating point cannot represent exactly, such as 0.1, Comet returns a wrong value where Spark returns NULL, 0.0 or 1.0 (https://github.com/apache/datafusion-comet/issues/6423) + +## RegrSXY + +The following incompatibilities cause `RegrSXY` to fall back to Spark by default. Set `spark.comet.expression.RegrSXY.allowIncompatible=true` to enable Comet acceleration despite these differences. + +- Comet merges the partial aggregates of `regr_sxy` in a different floating-point operation order from Spark. When a group's rows come from more than one partial aggregate and a variable is constant at a value that binary floating point cannot represent exactly, such as 0.1, Comet returns a wrong value where Spark returns NULL, 0.0 or 1.0 (https://github.com/apache/datafusion-comet/issues/6423) + +## RegrSlope + +The following incompatibilities cause `RegrSlope` to fall back to Spark by default. Set `spark.comet.expression.RegrSlope.allowIncompatible=true` to enable Comet acceleration despite these differences. + +- Comet merges the partial aggregates of `regr_slope` in a different floating-point operation order from Spark. When a group's rows come from more than one partial aggregate and a variable is constant at a value that binary floating point cannot represent exactly, such as 0.1, Comet returns a wrong value where Spark returns NULL, 0.0 or 1.0 (https://github.com/apache/datafusion-comet/issues/6423) + diff --git a/docs/source/user-guide/latest/compatibility/expressions/spark-3.4/array.md b/docs/source/user-guide/latest/compatibility/expressions/spark-3.4/array.md index 7bc6e299573..c723ab4e563 100644 --- a/docs/source/user-guide/latest/compatibility/expressions/spark-3.4/array.md +++ b/docs/source/user-guide/latest/compatibility/expressions/spark-3.4/array.md @@ -27,6 +27,12 @@ By default, `ArrayContains` is evaluated in the JVM using Spark's own code-gener - Spark compares array elements with ordering.equiv, so -0.0 matches +0.0 and all NaNs match each other; Comet's native array_contains compares the raw Arrow values bitwise +## ArrayDistinct + +The following incompatibilities cause `ArrayDistinct` to fall back to Spark by default. Set `spark.comet.expression.ArrayDistinct.allowIncompatible=true` to enable Comet acceleration despite these differences. + +- Floating-point elements match Spark's signed-zero and NaN semantics natively only on Spark 4.2.0, whose optimizer normalizes the arguments (SPARK-54918) + ## ArrayExcept By default, `ArrayExcept` is evaluated in the JVM using Spark's own code-generated implementation (run inside the Comet pipeline), which matches Spark exactly. Set `spark.comet.expression.ArrayExcept.allowIncompatible=true` to opt into Comet's native implementation instead, which has the following differences from Spark: @@ -47,6 +53,12 @@ By default, `ArrayJoin` is evaluated in the JVM using Spark's own code-generated - array_join does not propagate non-UTF8_BINARY collations to the output string (https://github.com/apache/datafusion-comet/issues/2190) - array_join evaluates its delimiter and null replacement eagerly, while Spark short-circuits past them (https://github.com/apache/datafusion-comet/issues/3178) +## ArrayUnion + +The following incompatibilities cause `ArrayUnion` to fall back to Spark by default. Set `spark.comet.expression.ArrayUnion.allowIncompatible=true` to enable Comet acceleration despite these differences. + +- Floating-point elements match Spark's signed-zero and NaN semantics natively only on Spark 4.2.0, whose optimizer normalizes the arguments (SPARK-54918) + ## ArraysZip The following cases are not supported by Comet and always fall back to Spark, regardless of any `allowIncompatible` setting: diff --git a/docs/source/user-guide/latest/compatibility/expressions/spark-3.5/aggregate.md b/docs/source/user-guide/latest/compatibility/expressions/spark-3.5/aggregate.md index daf7a9c00a5..22300066ce6 100644 --- a/docs/source/user-guide/latest/compatibility/expressions/spark-3.5/aggregate.md +++ b/docs/source/user-guide/latest/compatibility/expressions/spark-3.5/aggregate.md @@ -97,4 +97,34 @@ The following cases are not supported by Comet and always fall back to Spark, re - Descending order in `WITHIN GROUP (ORDER BY ... DESC)` is not supported. - Only numeric input types are supported. +## RegrIntercept + +The following incompatibilities cause `RegrIntercept` to fall back to Spark by default. Set `spark.comet.expression.RegrIntercept.allowIncompatible=true` to enable Comet acceleration despite these differences. + +- Comet merges the partial aggregates of `regr_intercept` in a different floating-point operation order from Spark. When a group's rows come from more than one partial aggregate and a variable is constant at a value that binary floating point cannot represent exactly, such as 0.1, Comet returns a wrong value where Spark returns NULL, 0.0 or 1.0 (https://github.com/apache/datafusion-comet/issues/6423) + +## RegrR2 + +The following incompatibilities cause `RegrR2` to fall back to Spark by default. Set `spark.comet.expression.RegrR2.allowIncompatible=true` to enable Comet acceleration despite these differences. + +- Comet merges the partial aggregates of `regr_r2` in a different floating-point operation order from Spark. When a group's rows come from more than one partial aggregate and a variable is constant at a value that binary floating point cannot represent exactly, such as 0.1, Comet returns a wrong value where Spark returns NULL, 0.0 or 1.0 (https://github.com/apache/datafusion-comet/issues/6423) + +## RegrReplacement + +The following incompatibilities cause `RegrReplacement` to fall back to Spark by default. Set `spark.comet.expression.RegrReplacement.allowIncompatible=true` to enable Comet acceleration despite these differences. + +- Comet merges the partial aggregates of `regr_sxx` and `regr_syy` in a different floating-point operation order from Spark. When a group's rows come from more than one partial aggregate and a variable is constant at a value that binary floating point cannot represent exactly, such as 0.1, Comet returns a wrong value where Spark returns NULL, 0.0 or 1.0 (https://github.com/apache/datafusion-comet/issues/6423) + +## RegrSXY + +The following incompatibilities cause `RegrSXY` to fall back to Spark by default. Set `spark.comet.expression.RegrSXY.allowIncompatible=true` to enable Comet acceleration despite these differences. + +- Comet merges the partial aggregates of `regr_sxy` in a different floating-point operation order from Spark. When a group's rows come from more than one partial aggregate and a variable is constant at a value that binary floating point cannot represent exactly, such as 0.1, Comet returns a wrong value where Spark returns NULL, 0.0 or 1.0 (https://github.com/apache/datafusion-comet/issues/6423) + +## RegrSlope + +The following incompatibilities cause `RegrSlope` to fall back to Spark by default. Set `spark.comet.expression.RegrSlope.allowIncompatible=true` to enable Comet acceleration despite these differences. + +- Comet merges the partial aggregates of `regr_slope` in a different floating-point operation order from Spark. When a group's rows come from more than one partial aggregate and a variable is constant at a value that binary floating point cannot represent exactly, such as 0.1, Comet returns a wrong value where Spark returns NULL, 0.0 or 1.0 (https://github.com/apache/datafusion-comet/issues/6423) + diff --git a/docs/source/user-guide/latest/compatibility/expressions/spark-3.5/array.md b/docs/source/user-guide/latest/compatibility/expressions/spark-3.5/array.md index 7bc6e299573..c723ab4e563 100644 --- a/docs/source/user-guide/latest/compatibility/expressions/spark-3.5/array.md +++ b/docs/source/user-guide/latest/compatibility/expressions/spark-3.5/array.md @@ -27,6 +27,12 @@ By default, `ArrayContains` is evaluated in the JVM using Spark's own code-gener - Spark compares array elements with ordering.equiv, so -0.0 matches +0.0 and all NaNs match each other; Comet's native array_contains compares the raw Arrow values bitwise +## ArrayDistinct + +The following incompatibilities cause `ArrayDistinct` to fall back to Spark by default. Set `spark.comet.expression.ArrayDistinct.allowIncompatible=true` to enable Comet acceleration despite these differences. + +- Floating-point elements match Spark's signed-zero and NaN semantics natively only on Spark 4.2.0, whose optimizer normalizes the arguments (SPARK-54918) + ## ArrayExcept By default, `ArrayExcept` is evaluated in the JVM using Spark's own code-generated implementation (run inside the Comet pipeline), which matches Spark exactly. Set `spark.comet.expression.ArrayExcept.allowIncompatible=true` to opt into Comet's native implementation instead, which has the following differences from Spark: @@ -47,6 +53,12 @@ By default, `ArrayJoin` is evaluated in the JVM using Spark's own code-generated - array_join does not propagate non-UTF8_BINARY collations to the output string (https://github.com/apache/datafusion-comet/issues/2190) - array_join evaluates its delimiter and null replacement eagerly, while Spark short-circuits past them (https://github.com/apache/datafusion-comet/issues/3178) +## ArrayUnion + +The following incompatibilities cause `ArrayUnion` to fall back to Spark by default. Set `spark.comet.expression.ArrayUnion.allowIncompatible=true` to enable Comet acceleration despite these differences. + +- Floating-point elements match Spark's signed-zero and NaN semantics natively only on Spark 4.2.0, whose optimizer normalizes the arguments (SPARK-54918) + ## ArraysZip The following cases are not supported by Comet and always fall back to Spark, regardless of any `allowIncompatible` setting: diff --git a/docs/source/user-guide/latest/compatibility/expressions/spark-4.0/aggregate.md b/docs/source/user-guide/latest/compatibility/expressions/spark-4.0/aggregate.md index 37fe120ae2a..3417acffc4c 100644 --- a/docs/source/user-guide/latest/compatibility/expressions/spark-4.0/aggregate.md +++ b/docs/source/user-guide/latest/compatibility/expressions/spark-4.0/aggregate.md @@ -106,4 +106,34 @@ The following cases are not supported by Comet and always fall back to Spark, re - Descending order in `WITHIN GROUP (ORDER BY ... DESC)` is not supported. - Only numeric input types are supported. +## RegrIntercept + +The following incompatibilities cause `RegrIntercept` to fall back to Spark by default. Set `spark.comet.expression.RegrIntercept.allowIncompatible=true` to enable Comet acceleration despite these differences. + +- Comet merges the partial aggregates of `regr_intercept` in a different floating-point operation order from Spark. When a group's rows come from more than one partial aggregate and a variable is constant at a value that binary floating point cannot represent exactly, such as 0.1, Comet returns a wrong value where Spark returns NULL, 0.0 or 1.0 (https://github.com/apache/datafusion-comet/issues/6423) + +## RegrR2 + +The following incompatibilities cause `RegrR2` to fall back to Spark by default. Set `spark.comet.expression.RegrR2.allowIncompatible=true` to enable Comet acceleration despite these differences. + +- Comet merges the partial aggregates of `regr_r2` in a different floating-point operation order from Spark. When a group's rows come from more than one partial aggregate and a variable is constant at a value that binary floating point cannot represent exactly, such as 0.1, Comet returns a wrong value where Spark returns NULL, 0.0 or 1.0 (https://github.com/apache/datafusion-comet/issues/6423) + +## RegrReplacement + +The following incompatibilities cause `RegrReplacement` to fall back to Spark by default. Set `spark.comet.expression.RegrReplacement.allowIncompatible=true` to enable Comet acceleration despite these differences. + +- Comet merges the partial aggregates of `regr_sxx` and `regr_syy` in a different floating-point operation order from Spark. When a group's rows come from more than one partial aggregate and a variable is constant at a value that binary floating point cannot represent exactly, such as 0.1, Comet returns a wrong value where Spark returns NULL, 0.0 or 1.0 (https://github.com/apache/datafusion-comet/issues/6423) + +## RegrSXY + +The following incompatibilities cause `RegrSXY` to fall back to Spark by default. Set `spark.comet.expression.RegrSXY.allowIncompatible=true` to enable Comet acceleration despite these differences. + +- Comet merges the partial aggregates of `regr_sxy` in a different floating-point operation order from Spark. When a group's rows come from more than one partial aggregate and a variable is constant at a value that binary floating point cannot represent exactly, such as 0.1, Comet returns a wrong value where Spark returns NULL, 0.0 or 1.0 (https://github.com/apache/datafusion-comet/issues/6423) + +## RegrSlope + +The following incompatibilities cause `RegrSlope` to fall back to Spark by default. Set `spark.comet.expression.RegrSlope.allowIncompatible=true` to enable Comet acceleration despite these differences. + +- Comet merges the partial aggregates of `regr_slope` in a different floating-point operation order from Spark. When a group's rows come from more than one partial aggregate and a variable is constant at a value that binary floating point cannot represent exactly, such as 0.1, Comet returns a wrong value where Spark returns NULL, 0.0 or 1.0 (https://github.com/apache/datafusion-comet/issues/6423) + diff --git a/docs/source/user-guide/latest/compatibility/expressions/spark-4.0/array.md b/docs/source/user-guide/latest/compatibility/expressions/spark-4.0/array.md index 7bc6e299573..c723ab4e563 100644 --- a/docs/source/user-guide/latest/compatibility/expressions/spark-4.0/array.md +++ b/docs/source/user-guide/latest/compatibility/expressions/spark-4.0/array.md @@ -27,6 +27,12 @@ By default, `ArrayContains` is evaluated in the JVM using Spark's own code-gener - Spark compares array elements with ordering.equiv, so -0.0 matches +0.0 and all NaNs match each other; Comet's native array_contains compares the raw Arrow values bitwise +## ArrayDistinct + +The following incompatibilities cause `ArrayDistinct` to fall back to Spark by default. Set `spark.comet.expression.ArrayDistinct.allowIncompatible=true` to enable Comet acceleration despite these differences. + +- Floating-point elements match Spark's signed-zero and NaN semantics natively only on Spark 4.2.0, whose optimizer normalizes the arguments (SPARK-54918) + ## ArrayExcept By default, `ArrayExcept` is evaluated in the JVM using Spark's own code-generated implementation (run inside the Comet pipeline), which matches Spark exactly. Set `spark.comet.expression.ArrayExcept.allowIncompatible=true` to opt into Comet's native implementation instead, which has the following differences from Spark: @@ -47,6 +53,12 @@ By default, `ArrayJoin` is evaluated in the JVM using Spark's own code-generated - array_join does not propagate non-UTF8_BINARY collations to the output string (https://github.com/apache/datafusion-comet/issues/2190) - array_join evaluates its delimiter and null replacement eagerly, while Spark short-circuits past them (https://github.com/apache/datafusion-comet/issues/3178) +## ArrayUnion + +The following incompatibilities cause `ArrayUnion` to fall back to Spark by default. Set `spark.comet.expression.ArrayUnion.allowIncompatible=true` to enable Comet acceleration despite these differences. + +- Floating-point elements match Spark's signed-zero and NaN semantics natively only on Spark 4.2.0, whose optimizer normalizes the arguments (SPARK-54918) + ## ArraysZip The following cases are not supported by Comet and always fall back to Spark, regardless of any `allowIncompatible` setting: diff --git a/docs/source/user-guide/latest/compatibility/expressions/spark-4.1/aggregate.md b/docs/source/user-guide/latest/compatibility/expressions/spark-4.1/aggregate.md index 37fe120ae2a..3417acffc4c 100644 --- a/docs/source/user-guide/latest/compatibility/expressions/spark-4.1/aggregate.md +++ b/docs/source/user-guide/latest/compatibility/expressions/spark-4.1/aggregate.md @@ -106,4 +106,34 @@ The following cases are not supported by Comet and always fall back to Spark, re - Descending order in `WITHIN GROUP (ORDER BY ... DESC)` is not supported. - Only numeric input types are supported. +## RegrIntercept + +The following incompatibilities cause `RegrIntercept` to fall back to Spark by default. Set `spark.comet.expression.RegrIntercept.allowIncompatible=true` to enable Comet acceleration despite these differences. + +- Comet merges the partial aggregates of `regr_intercept` in a different floating-point operation order from Spark. When a group's rows come from more than one partial aggregate and a variable is constant at a value that binary floating point cannot represent exactly, such as 0.1, Comet returns a wrong value where Spark returns NULL, 0.0 or 1.0 (https://github.com/apache/datafusion-comet/issues/6423) + +## RegrR2 + +The following incompatibilities cause `RegrR2` to fall back to Spark by default. Set `spark.comet.expression.RegrR2.allowIncompatible=true` to enable Comet acceleration despite these differences. + +- Comet merges the partial aggregates of `regr_r2` in a different floating-point operation order from Spark. When a group's rows come from more than one partial aggregate and a variable is constant at a value that binary floating point cannot represent exactly, such as 0.1, Comet returns a wrong value where Spark returns NULL, 0.0 or 1.0 (https://github.com/apache/datafusion-comet/issues/6423) + +## RegrReplacement + +The following incompatibilities cause `RegrReplacement` to fall back to Spark by default. Set `spark.comet.expression.RegrReplacement.allowIncompatible=true` to enable Comet acceleration despite these differences. + +- Comet merges the partial aggregates of `regr_sxx` and `regr_syy` in a different floating-point operation order from Spark. When a group's rows come from more than one partial aggregate and a variable is constant at a value that binary floating point cannot represent exactly, such as 0.1, Comet returns a wrong value where Spark returns NULL, 0.0 or 1.0 (https://github.com/apache/datafusion-comet/issues/6423) + +## RegrSXY + +The following incompatibilities cause `RegrSXY` to fall back to Spark by default. Set `spark.comet.expression.RegrSXY.allowIncompatible=true` to enable Comet acceleration despite these differences. + +- Comet merges the partial aggregates of `regr_sxy` in a different floating-point operation order from Spark. When a group's rows come from more than one partial aggregate and a variable is constant at a value that binary floating point cannot represent exactly, such as 0.1, Comet returns a wrong value where Spark returns NULL, 0.0 or 1.0 (https://github.com/apache/datafusion-comet/issues/6423) + +## RegrSlope + +The following incompatibilities cause `RegrSlope` to fall back to Spark by default. Set `spark.comet.expression.RegrSlope.allowIncompatible=true` to enable Comet acceleration despite these differences. + +- Comet merges the partial aggregates of `regr_slope` in a different floating-point operation order from Spark. When a group's rows come from more than one partial aggregate and a variable is constant at a value that binary floating point cannot represent exactly, such as 0.1, Comet returns a wrong value where Spark returns NULL, 0.0 or 1.0 (https://github.com/apache/datafusion-comet/issues/6423) + diff --git a/docs/source/user-guide/latest/compatibility/expressions/spark-4.1/array.md b/docs/source/user-guide/latest/compatibility/expressions/spark-4.1/array.md index 7bc6e299573..c723ab4e563 100644 --- a/docs/source/user-guide/latest/compatibility/expressions/spark-4.1/array.md +++ b/docs/source/user-guide/latest/compatibility/expressions/spark-4.1/array.md @@ -27,6 +27,12 @@ By default, `ArrayContains` is evaluated in the JVM using Spark's own code-gener - Spark compares array elements with ordering.equiv, so -0.0 matches +0.0 and all NaNs match each other; Comet's native array_contains compares the raw Arrow values bitwise +## ArrayDistinct + +The following incompatibilities cause `ArrayDistinct` to fall back to Spark by default. Set `spark.comet.expression.ArrayDistinct.allowIncompatible=true` to enable Comet acceleration despite these differences. + +- Floating-point elements match Spark's signed-zero and NaN semantics natively only on Spark 4.2.0, whose optimizer normalizes the arguments (SPARK-54918) + ## ArrayExcept By default, `ArrayExcept` is evaluated in the JVM using Spark's own code-generated implementation (run inside the Comet pipeline), which matches Spark exactly. Set `spark.comet.expression.ArrayExcept.allowIncompatible=true` to opt into Comet's native implementation instead, which has the following differences from Spark: @@ -47,6 +53,12 @@ By default, `ArrayJoin` is evaluated in the JVM using Spark's own code-generated - array_join does not propagate non-UTF8_BINARY collations to the output string (https://github.com/apache/datafusion-comet/issues/2190) - array_join evaluates its delimiter and null replacement eagerly, while Spark short-circuits past them (https://github.com/apache/datafusion-comet/issues/3178) +## ArrayUnion + +The following incompatibilities cause `ArrayUnion` to fall back to Spark by default. Set `spark.comet.expression.ArrayUnion.allowIncompatible=true` to enable Comet acceleration despite these differences. + +- Floating-point elements match Spark's signed-zero and NaN semantics natively only on Spark 4.2.0, whose optimizer normalizes the arguments (SPARK-54918) + ## ArraysZip The following cases are not supported by Comet and always fall back to Spark, regardless of any `allowIncompatible` setting: diff --git a/docs/source/user-guide/latest/configs.md b/docs/source/user-guide/latest/configs.md index 035b074593a..c6981536c4d 100644 --- a/docs/source/user-guide/latest/configs.md +++ b/docs/source/user-guide/latest/configs.md @@ -58,6 +58,7 @@ Comet provides the following configuration settings. | `spark.comet.convert.parquet.enabled` | When enabled, data from Spark (non-native) Parquet v1 and v2 scans will be converted to Arrow format. | false | | `spark.comet.debug.enabled` | Whether to enable debug mode for Comet. When enabled, Comet will do additional checks for debugging purpose. For example, validating array when importing arrays from JVM at native side. Note that these checks may be expensive in performance and should only be enabled for debugging purpose. | false | | `spark.comet.enabled` | Whether to enable Comet extension for Spark. When this is turned on, Spark will use Comet to read Parquet data source. Note that to enable native vectorized execution, both this config and `spark.comet.exec.enabled` need to be enabled. It can be overridden by the environment variable `ENABLE_COMET`. | true | +| `spark.comet.exec.aggregate.skipPartial.enabled` | Experimental opt-in: let a native partial aggregate stop aggregating once its input looks mostly distinct, and send the rest of the task's rows to the shuffle unaggregated. Only applies to partial aggregates that feed Comet native shuffle and whose aggregate functions, if any, are all single-argument COUNT. The check starts after the first 100,000 input rows of a task, and once aggregation stops it does not resume, so a task whose keys repeat after a mostly distinct start can shuffle many times more rows than it would with this disabled. For more information, refer to the [Comet Tuning Guide](https://datafusion.apache.org/comet/user-guide/latest/tuning.html). | false | | `spark.comet.exec.columnarToRow.native.enabled` | Whether to enable native columnar to row conversion. When enabled, Comet will use native Rust code to convert Arrow columnar data to Spark UnsafeRow format instead of the JVM implementation. The native conversion carries a fixed JNI cost per batch and is slower than the JVM implementation for small batches. | false | | `spark.comet.exec.enabled` | Whether to enable Comet native vectorized execution for Spark. This controls whether Spark should convert operators into their Comet counterparts and execute them in native space. Note: each operator is associated with a separate config in the format of `spark.comet.exec..enabled` at the moment, and both the config and this need to be turned on, in order for the operator to be executed in native. | true | | `spark.comet.exec.forceShuffledHashJoin` | Experimental feature to force Spark to replace SortMergeJoin with ShuffledHashJoin for improved performance. This feature is not stable yet. For more information, refer to the [Comet Tuning Guide](https://datafusion.apache.org/comet/user-guide/latest/tuning.html). | false | diff --git a/docs/source/user-guide/latest/expressions.md b/docs/source/user-guide/latest/expressions.md index f010a642332..d0f560489e7 100644 --- a/docs/source/user-guide/latest/expressions.md +++ b/docs/source/user-guide/latest/expressions.md @@ -125,9 +125,9 @@ The tables below list every Spark built-in expression with its current status. | `regr_intercept` | ✅ | Native | Falls back by default because the native merge of partial aggregates differs from Spark ([#6423](https://github.com/apache/datafusion-comet/issues/6423)); the native path is opt-in via `spark.comet.expression.RegrIntercept.allowIncompatible=true` | | `regr_r2` | ✅ | Native | Falls back by default because the native merge of partial aggregates differs from Spark ([#6423](https://github.com/apache/datafusion-comet/issues/6423)); the native path is opt-in via `spark.comet.expression.RegrR2.allowIncompatible=true` | | `regr_slope` | ✅ | Native | Falls back by default because the native merge of partial aggregates differs from Spark ([#6423](https://github.com/apache/datafusion-comet/issues/6423)); the native path is opt-in via `spark.comet.expression.RegrSlope.allowIncompatible=true` | -| `regr_sxx` | ✅ | Native | Falls back by default because the native merge of partial aggregates differs from Spark ([#6423](https://github.com/apache/datafusion-comet/issues/6423)); the native path is opt-in via `spark.comet.expression.RegrReplacement.allowIncompatible=true` (Spark plans `regr_sxx` as `RegrReplacement`) | +| `regr_sxx` | ✅ | — | Falls back by default because the native merge of partial aggregates differs from Spark ([#6423](https://github.com/apache/datafusion-comet/issues/6423)); the native path is opt-in via `spark.comet.expression.RegrReplacement.allowIncompatible=true` (Spark plans `regr_sxx` as `RegrReplacement`) | | `regr_sxy` | ✅ | Native | Falls back by default because the native merge of partial aggregates differs from Spark ([#6423](https://github.com/apache/datafusion-comet/issues/6423)); the native path is opt-in via `spark.comet.expression.RegrSXY.allowIncompatible=true` | -| `regr_syy` | ✅ | Native | Falls back by default because the native merge of partial aggregates differs from Spark ([#6423](https://github.com/apache/datafusion-comet/issues/6423)); the native path is opt-in via `spark.comet.expression.RegrReplacement.allowIncompatible=true` (Spark plans `regr_syy` as `RegrReplacement`) | +| `regr_syy` | ✅ | — | Falls back by default because the native merge of partial aggregates differs from Spark ([#6423](https://github.com/apache/datafusion-comet/issues/6423)); the native path is opt-in via `spark.comet.expression.RegrReplacement.allowIncompatible=true` (Spark plans `regr_syy` as `RegrReplacement`) | | `skewness` | 🔜 | — | Not yet implemented natively | | `some` | ✅ | — | | | `std` | ✅ | Native | |