andygrove commented on code in PR #5634: URL: https://github.com/apache/datafusion-comet/pull/5634#discussion_r4232775813
########## docs/source/user-guide/latest/in-memory-cache.md: ########## @@ -21,18 +21,24 @@ Comet can store Spark's in-memory cache (`CACHE TABLE`, `df.cache()`, `df.persist()`) in an Arrow format that Comet operators read directly. Without it, a cached table is stored in Spark's own -format and every scan of it has to convert each batch before Comet can continue, which shows up in -the plan as a `CometSparkColumnarToColumnar` above the cache scan. +format, which Comet operators cannot read. Under Comet's default settings the operators above the +cache scan then run on Spark. With `spark.comet.convert.inMemoryCache.enabled`, a +`CometSparkColumnarToColumnar` above the scan converts each batch for Comet operators instead. -This feature is **experimental and disabled by default**. Turn it on at startup, alongside the rest -of Comet's configuration: +This feature is **experimental and enabled by default** from Spark 3.5. To turn it off, set the Review Comment: Agreed, it shouldn't be both. 684c1e0681 drops "experimental" here, on the compatibility page and on the operators page, which now marks `InMemoryTableScanExec` ✅. It also fixes the plan-node table in Understanding Comet Plans, which still said the cache scan was disabled by default. ########## docs/source/user-guide/latest/operators.md: ########## @@ -66,7 +66,7 @@ omitted from the tables below and may be reconsidered based on demand: | `LocalTableScanExec` | ⚠️ | Disabled by default; there is no acceleration advantage and this operator is typically only used in test code. Can be opted into via config ([#4393](https://github.com/apache/datafusion-comet/pull/4393)). | | `EmptyRelationExec` | ✅ | Spark 4.0 and later. See [Empty Relations](compatibility/operators.md#empty-relations) for native-input support and writer fallback. | | `RangeExec` | ⚠️ | Disabled by default. Set `spark.comet.exec.range.enabled=true` to generate the rows of `spark.range` and SQL `range()` in native code, so the operators above them run natively. It can be slower than Spark when those operators are only cheap expressions, such as a filter, which Spark compiles together with the range into one loop. Ranges whose arithmetic overflows the `Long` range fall back to Spark. | -| `InMemoryTableScanExec` | ⚠️ | Experimental, disabled by default. Set `spark.comet.exec.inMemoryCache.enabled=true` before the application starts so Comet installs its Arrow cache serializer. Relations with unsupported column types stay in Spark's cache format and fall back. See [In-Memory Cache](in-memory-cache.md). | +| `InMemoryTableScanExec` | ⚠️ | Experimental, enabled by default from Spark 3.5. Comet installs its Arrow cache serializer as the application starts if `spark.comet.exec.inMemoryCache.enabled` is true, its default except on Spark 3.4. Relations with unsupported column types stay in Spark's cache format and fall back. See [In-Memory Cache](in-memory-cache.md). | Review Comment: Removed in 684c1e0681. The row is ✅ now, and its note covers Spark 3.4, where the cache stays off by default. ########## dev/diffs/4.1.3.diff: ########## @@ -423,28 +474,63 @@ index 0d807aeae4d..6d7744e771b 100644 } test("A cached table preserves the partitioning and ordering of its cached SparkPlan") { -@@ -1673,9 +1674,18 @@ class CachedTableSuite extends QueryTest with SQLTestUtils +@@ -1595,7 +1602,8 @@ class CachedTableSuite extends QueryTest with SQLTestUtils + } + } + +- test("SPARK-36120: Support cache/uncache table with TimestampNTZ type") { ++ test("SPARK-36120: Support cache/uncache table with TimestampNTZ type", ++ IgnoreComet("Comet's cache format reports Arrow buffer sizes")) { Review Comment: Done in b4befe74a4, in the 3.5, 4.0, 4.1 and 4.2 diffs. The four tests now skip only the `sizeInBytes` checks in SPARK-36120 and SPARK-22673, the three `checkOptimizedPlanStats(..., 4, ...)` calls in SPARK-33687, and the reuse count in SPARK-43376. I ran them locally with Comet's format against each version's test jar, with the patched harness compiled in, and they pass on all four. Comet reports 9, 5 and 18 bytes where Spark expects 8, 4 and 16, and SPARK-43376 finds no reused subquery, while its answer check and the rest of `runAdaptiveAndVerifyResult` pass. ########## dev/diffs/4.1.3.diff: ########## @@ -4410,6 +4874,11 @@ index 720b13b812e..e3ac2cebc6e 100644 + "org.apache.spark.sql.comet.execution.shuffle.CometShuffleManager") + .set("spark.comet.shuffle.enabled", "true") + ++ // CometDriverPlugin installs Comet's cache serializer when ++ // spark.comet.exec.inMemoryCache.enabled is on, as it is by default. These sessions do not ++ // load the plugin, so install it here. ++ conf.set("spark.sql.cache.serializer", Review Comment: Agreed. b4befe74a4 adds it to `CometTestSettings` next to the shuffle manager, in the 3.5, 4.0, 4.1 and 4.2 diffs. The 3.4 diff stays as on `main`, since the default is off there. I checked the case you describe on 4.1: in a fresh JVM, `BroadcastJoinSuite`, `BroadcastJoinSuiteAE` and `SparkSessionExtensionSuite` cached in Spark's format without the property and in Comet's with it, and passed both ways. With it, they also pass on 3.5, 4.0 and 4.2. `SharedSparkSession` and `TestHive` still set the serializer themselves, for runs that don't go through sbt. ########## spark/src/test/scala/org/apache/spark/sql/benchmark/CometInMemoryCacheBenchmark.scala: ########## @@ -488,6 +485,116 @@ object CometInMemoryCacheBenchmark extends CometBenchmarkBase { } } + /** + * What the feature changes for a query that runs with Comet, against Spark's own cache format, + * with AQE on and Comet's other settings at their defaults. This is the comparison an + * application gets from turning the feature on or off. Both formats are cached from the same + * relation, one copy at a time as in runSparkOperatorBenchmark. + * + * Two shapes of plan read the cache. With Comet operators above the cache scan, Comet's format + * runs the whole query natively, while Spark's leaves the operators directly above its scan on + * Spark: Comet reads Spark's cache scan only through spark.comet.convert.inMemoryCache.enabled, + * which is off by default. With a Spark operator above the scan, Comet's format is read by the + * native scan and converted to rows for that operator, where Spark's is read by Spark's own + * scan. The Spark operator is the aggregate, with Comet's turned off, standing in for any + * operator Comet does not support. + */ + private def runAdaptiveBenchmark(relation: CachedRelation): Unit = { Review Comment: Added in 5c9cc74a15. `runBuildBenchmark` times building the cache in Spark's format and in Comet's with both codecs, from Comet operators and from Spark operators, and prints each footprint. On an M3 Max, Comet's format with `zstd` built in 944 ms from Comet operators and 2079 ms from Spark operators, against 3285 ms and 3680 ms for Spark's format, and it holds the relation in 51 to 55 MiB against Spark's 217 MiB (315 MiB with `none`). So the build is faster in Comet's format too: 1.8x when Spark operators produce the rows, and 3.5x when Comet operators do, since Spark's format then has to convert their batches to rows first. The table is in the cache guide, under "Against Spark's cache format". ########## dev/diffs/4.1.3.diff: ########## Review Comment: Agreed, and thanks for working out why. b4befe74a4 drops the guard from the 3.5, 4.0, 4.1 and 4.2 diffs and keeps it on 3.4. The test passes locally with Comet on all four versions, with a `CometUnion` whose two join-side reads are both coalesced. -- 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]
