Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
73 changes: 71 additions & 2 deletions docs/source/user-guide/latest/migration-guide.md
Original file line number Diff line number Diff line change
Expand Up @@ -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:

Expand Down Expand Up @@ -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
Expand Down
14 changes: 7 additions & 7 deletions spark/src/main/scala/org/apache/comet/CometConf.scala
Original file line number Diff line number Diff line change
Expand Up @@ -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)

Expand Down Expand Up @@ -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)
Expand Down
8 changes: 4 additions & 4 deletions spark/src/main/scala/org/apache/spark/Plugins.scala
Original file line number Diff line number Diff line change
Expand Up @@ -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}.")
}
}
Expand Down
Loading