diff --git a/docs/source/index.md b/docs/source/index.md index 5bcaedea8b4..0eec6635207 100644 --- a/docs/source/index.md +++ b/docs/source/index.md @@ -139,6 +139,7 @@ Runs your existing Spark queries on the Apache DataFusion native engine, no code User Guide Contributor Guide +Release Notes Changelog About ``` diff --git a/docs/source/release-notes/1.1.0.md b/docs/source/release-notes/1.1.0.md new file mode 100644 index 00000000000..939feb13397 --- /dev/null +++ b/docs/source/release-notes/1.1.0.md @@ -0,0 +1,56 @@ + + +# Comet 1.1.0 Release Notes + +Comet 1.1.0 fixes about 90 bugs that shipped in 1.0.0, about 40 of which returned results that +differed from Spark. Before upgrading, read the +[upgrade guide](../user-guide/latest/migration-guide.md#upgrading-to-comet-110) for the behavior +changes in this release, and the known regressions below. + +## Known Regressions + +After the first release candidate was cut, an audit of every pull request in 1.1.0 against 1.0.0 +looked for cases that 1.0.0 handled correctly and 1.1.0 doesn't +([#6399](https://github.com/apache/datafusion-comet/issues/6399)). This release fixes every +regression it found except the one below. Its entry says how to avoid the problem, and each setting +it names has been checked against the regression's reproducer on 1.1.0. A fix is planned for +1.1.1, and [#6402](https://github.com/apache/datafusion-comet/issues/6402) has the full list, +including the regressions that were fixed before the release. + +### Errors + +- **A spilled native aggregate under memory pressure** + ([#6254](https://github.com/apache/datafusion-comet/issues/6254)). After a native final hash + aggregate spills, it reads the spilled data back with no way to spill again. 1.0.0 ignored a + refused memory request at that step, but since the upgrade to DataFusion 55 the refusal fails the + task. So a task near its memory limit can fail with `Failed to acquire N bytes` in + `FinalHashAggregateStream`. To avoid it, give executors more off-heap memory with + `spark.memory.offHeap.size`. Twice the size that failed was enough in our tests, but the margin + depends on the workload. Setting `spark.comet.exec.aggregate.enabled=false` for the affected job + also avoids it, by running its aggregates in Spark. + +### Performance + +- **Partial aggregates without the Comet shuffle manager.** When the final aggregate runs in Spark, + which happens for every aggregate if the Comet shuffle manager isn't installed, the partial `avg`, + decimal `sum`, `stddev`, `variance`, `corr`, `first`, `last` and a few others now run in Spark + too. In 1.0.0 they ran natively, and `avg` could return NULL in this plan + ([#5419](https://github.com/apache/datafusion-comet/issues/5419)). Installing the Comet shuffle + manager keeps them native. diff --git a/docs/source/release-notes/index.md b/docs/source/release-notes/index.md new file mode 100644 index 00000000000..c915633f36c --- /dev/null +++ b/docs/source/release-notes/index.md @@ -0,0 +1,31 @@ + + +# Release Notes + +The release notes describe what users of each Comet release should know that the change log and the +[upgrade guide](../user-guide/latest/migration-guide.md) don't cover, such as known regressions and +the settings that avoid them. The [change log](../changelog/index.md) lists every pull request in a +release. + +```{toctree} +:maxdepth: 1 + +1.1.0 <1.1.0> +``` diff --git a/docs/source/user-guide/latest/migration-guide.md b/docs/source/user-guide/latest/migration-guide.md index 342f8d21182..81421af9c7f 100644 --- a/docs/source/user-guide/latest/migration-guide.md +++ b/docs/source/user-guide/latest/migration-guide.md @@ -62,6 +62,9 @@ key is removed. 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` also has a known regression, with a setting that avoids it. It is described in the +[1.1.0 release notes](../../release-notes/1.1.0.md). + 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. @@ -124,6 +127,19 @@ session's `spark.shuffle.manager`. A session that named `CometShuffleManager` af 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. +### Native Iceberg Reads on EKS with IRSA + +On EKS with IAM Roles for Service Accounts (IRSA), when `AWS_WEB_IDENTITY_TOKEN_FILE`, +`AWS_ROLE_ARN` and a region are set and the catalog configures no credentials, Comet `1.1.0`'s +native Iceberg scan takes its S3 credentials only from the web-identity role. If that fails, it no +longer falls back to the node role or Pod Identity, as Comet `1.0.0` did. So a cluster whose IRSA +setup is broken, and that was reading S3 as the node role without anyone noticing, now fails native +Iceberg reads with `failed to load signing credential`. The "EKS / IRSA" section of +[S3 Credential Providers](s3-credential-providers.md) explains the change. To go back to the old +credential chain for a catalog, set +`spark.sql.catalog..s3.comet.credential.webIdentity.enabled=false`. A table loaded by path +has no catalog to set that on, so for it set `spark.comet.scan.icebergNative.enabled=false`. + ### Deprecated and Removed Settings `spark.comet.exec.memoryPool.fraction` is deprecated and will be removed in a future major release.