diff --git a/docs/source/user-guide/latest/migration-guide.md b/docs/source/user-guide/latest/migration-guide.md index 1de5f88467..a5b00dbc23 100644 --- a/docs/source/user-guide/latest/migration-guide.md +++ b/docs/source/user-guide/latest/migration-guide.md @@ -27,7 +27,11 @@ A **behavior change** is one where the same query, run over the same data, with set configuration, produces a different result or a different error than it did in the previous release. Comet's [versioning policy](../../about/versioning_policy.md) permits these in a minor release only when a `spark.comet.legacy.*` configuration key restores the previous behavior, so -every entry below names such a key. +every behavior change below names such a key. + +A release's section can also list changes that need no legacy key but can still change what an +existing deployment does, such as a fix that makes Comet apply a setting as documented, or a setting +that was deprecated or removed. Each of these entries names the settings involved. Two kinds of change are deliberately absent from this guide: @@ -55,11 +59,76 @@ key is removed. ## Upgrading to Comet 1.1.0 -Comet `1.1.0` makes no behavior changes that need a `spark.comet.legacy.*` key. +Comet `1.1.0` makes no behavior changes that need a `spark.comet.legacy.*` key. The changes below +need none either, but check whether any of them applies to your deployment. Comet `1.1.0` requires JDK 17 or later. JDK 11 is no longer supported. See [Installing Comet](installation.md) for the supported Java, Scala, and Spark versions. +### Settings That Now Take Effect as Documented + +Comet `1.0.0` misread three size settings. Comet `1.1.0` reads each of them as documented, so a job +that sets one of them can behave differently after the upgrade. In each case, a setting that +already exists restores the old effect. + +- `spark.comet.shuffle.native.writeBufferSize` was read in MiB but used as a number of bytes, so + the native shuffle writer ran with a 1-byte write buffer by default, and a value of `64m` gave it + 64 bytes. The setting is now read in bytes with a default of 1 MiB, and a value with a unit means + what it says. Each shuffle task holds a few of these buffers in native memory that no memory pool + tracks, so check any value you set with a unit against `spark.executor.memoryOverhead`. A bare + number keeps its old meaning, so writing the number without its unit restores the old buffer + size. Values of 2 GiB or more are now rejected. +- `spark.comet.maxTempDirectorySize` was ignored when it was written with a unit, such as `10g`, and + the 100 GB default applied instead. It is now enforced, so a query that spills more than that + amount now fails. Remove the setting to keep the old limit. See + [Limiting Spill Disk Usage](tuning/memory.md#limiting-spill-disk-usage). +- `spark.memory.offHeap.size` was read as MiB when it was written as a bare number of bytes. The + `fair_unified` memory pool's per-operator shares were therefore about a million times too large, + and never limited an operator. The shares are now correct, so operators can spill sooner, and an + operator that cannot spill can fail when it exceeds its share. A size written with a unit, such + as `16g`, is unaffected. To get the old behavior back, set + `spark.comet.exec.memoryPool=greedy_unified`, which leaves every limit to Spark. See + [Configuring Comet Memory](tuning/memory.md#configuring-comet-memory). + +A malformed value of `spark.comet.maxTempDirectorySize` or `spark.comet.explain.native.enabled` now +fails the query instead of being replaced by the default. `spark.comet.debug.enabled`, +`spark.comet.explain.native.enabled` and `spark.comet.tracing.enabled` now also take effect in +Comet's native code when they are written in upper case, such as `TRUE`. + +### Conditions for Enabling Comet + +Comet needs Spark's off-heap memory to be enabled. `CometPlugin` already disabled Comet when +off-heap memory was disabled, but an application that registered `CometSparkSessionExtensions` +directly with `spark.sql.extensions` skipped that check, and Comet ran in on-heap mode, which exists +only for running tests. Comet `1.1.0` makes the same check on that path, and disables itself with a +warning when off-heap memory is not enabled. To keep using Comet, set +`spark.memory.offHeap.enabled=true` and `spark.memory.offHeap.size` when the application starts; +see [Configuring Comet Memory](tuning/memory.md#configuring-comet-memory). Both checks read +`spark.memory.offHeap.enabled` from the SparkContext, so setting it on a `SparkSession.builder` +after the SparkContext exists has no effect. + +Comet also now checks the shuffle manager that the application is running, rather than the +session's `spark.shuffle.manager`. A session that named `CometShuffleManager` after the SparkContext +had started with a different shuffle manager used to plan Comet shuffles that failed with a +`ClassCastException`. Such a session now runs without Comet, with a warning. + +### Deprecated and Removed Settings + +`spark.comet.exec.memoryPool.fraction` is deprecated and will be removed in a future major release. +It was documented as leaving room in `spark.memory.offHeap.size` for the native memory that Comet's +memory pools do not track, but it cannot: Spark hands out the whole off-heap pool whatever it is set +to. It keeps working as before, and the driver now logs a warning when it is set. Size +`spark.executor.memoryOverhead` for that memory instead; see +[Configuring Executor Memory Overhead](tuning/memory.md#configuring-executor-memory-overhead). + +Comet `1.1.0` removes `spark.comet.memoryOverhead`, `spark.comet.exec.onHeap.memoryPool` and +`spark.comet.shuffle.jvm.memoryFactor`, including its older name +`spark.comet.columnar.shuffle.memory.factor`. They applied only to on-heap mode +(`spark.comet.exec.onHeap.enabled`), which exists for running Spark's SQL tests against Comet and no +longer tracks native memory at all. They were in the testing category, which the +[versioning policy](../../about/versioning_policy.md#testing-and-internal-configurations-are-exempt) +exempts, and Comet ignores them if they are still set. + ## Upgrading to Comet 1.0.0 Comet `1.0.0` is the first release under the stable diff --git a/spark/src/main/scala/org/apache/comet/CometConf.scala b/spark/src/main/scala/org/apache/comet/CometConf.scala index eeb5cf8045..53daeb5710 100644 --- a/spark/src/main/scala/org/apache/comet/CometConf.scala +++ b/spark/src/main/scala/org/apache/comet/CometConf.scala @@ -336,7 +336,7 @@ object CometConf extends ShimCometConf { conf("spark.comet.exec.sortMergeJoinWithJoinFilter.enabled") .category(CATEGORY_ENABLE_EXEC) .doc("Support for Sort Merge Join with filter. " + - "Deprecated: this config will be removed in a future release.") + "Deprecated: this config will be removed in a future major release.") .booleanConf .createWithDefault(true) @@ -916,12 +916,12 @@ object CometConf extends ShimCometConf { conf("spark.comet.exec.memoryPool.fraction") .category(CATEGORY_TUNING) .doc( - "Deprecated: this config will be removed in a future release. It does not leave room " + - "in spark.memory.offHeap.size for native memory that Comet's memory pools do not " + - "track, because Spark hands out the whole off-heap pool whatever this is set to. Size " + - "spark.executor.memoryOverhead for that memory instead. Only applies to off-heap " + - "mode, where the fair_unified pool limits each memory consumer in a task to this " + - "fraction of the off-heap size divided by the task's consumers, and the " + + "Deprecated: this config will be removed in a future major release. It does not " + + "leave room in spark.memory.offHeap.size for native memory that Comet's memory " + + "pools do not track, because Spark hands out the whole off-heap pool whatever this " + + "is set to. Size spark.executor.memoryOverhead for that memory instead. Only applies " + + "to off-heap mode, where the fair_unified pool limits each memory consumer in a task " + + "to this fraction of the off-heap size divided by the task's consumers, and the " + s"greedy_unified pool ignores it. $TUNING_GUIDE.") .doubleConf .createWithDefault(1.0) diff --git a/spark/src/main/scala/org/apache/spark/Plugins.scala b/spark/src/main/scala/org/apache/spark/Plugins.scala index 2fe6cceff2..b0680ae695 100644 --- a/spark/src/main/scala/org/apache/spark/Plugins.scala +++ b/spark/src/main/scala/org/apache/spark/Plugins.scala @@ -215,10 +215,10 @@ object CometDriverPlugin extends Logging { val key = CometConf.COMET_OFFHEAP_MEMORY_POOL_FRACTION.key conf.getOption(key).foreach { value => logWarning( - s"$key=$value is deprecated and will be removed in a future release. It does not leave " + - "room in spark.memory.offHeap.size for native memory that Comet's memory pools do " + - "not track, because Spark hands out the whole off-heap pool whatever it is set to. " + - s"Size ${EXECUTOR_MEMORY_OVERHEAD.key} for that memory instead. " + + s"$key=$value is deprecated and will be removed in a future major release. It does " + + "not leave room in spark.memory.offHeap.size for native memory that Comet's memory " + + "pools do not track, because Spark hands out the whole off-heap pool whatever it is " + + s"set to. Size ${EXECUTOR_MEMORY_OVERHEAD.key} for that memory instead. " + s"${CometConf.TUNING_GUIDE}.") } }