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]

Reply via email to