andygrove commented on code in PR #5859:
URL: https://github.com/apache/datafusion-comet/pull/5859#discussion_r4128052535


##########
docs/source/user-guide/latest/in-memory-cache.md:
##########
@@ -175,19 +176,45 @@ registrator.
 
 ## Limitations
 
-Reads that feed **Spark** operators rather than Comet ones are slower than 
Spark's own cache
-format, and the narrower the read, the wider the gap. Measured by the same 
benchmark over the same
-5M-row relation, with Comet off so that Spark operators consume the cached 
data:
-
-| Read shape              | Spark's cache format | Comet's cache format | 
Slowdown |
-| ----------------------- | -------------------: | -------------------: | 
-------: |
-| Row count only (0 of 6) |                35 ms |               183 ms |     
5.2x |
-| 1 of 6 columns          |                54 ms |               257 ms |     
4.8x |
-| 3 of 6 columns          |                98 ms |               331 ms |     
3.4x |
-| 6 of 6 columns          |               410 ms |               623 ms |     
1.5x |
-
-This is why the feature is off by default. The cause is not yet established;
-[#5485](https://github.com/apache/datafusion-comet/issues/5485) tracks it.
+Reads that feed **Spark** operators rather than Comet ones can be slower than 
Spark's own cache
+format, because every cached batch is decoded from Arrow before Spark reads 
it. How a Spark
+operator reads a relation cached in Comet's format depends on the scan below 
it:
+
+- With native execution enabled (`spark.comet.exec.enabled=true`), the cache 
is scanned by
+  `CometInMemoryTableScan`, and a Spark operator above it reads the scan's 
batches through
+  `CometColumnarToRow`, as it would above any other Comet operator.
+- With Comet enabled but native execution disabled, the cache is scanned by 
Spark's
+  `InMemoryTableScanExec`. When the operator directly above the scan takes 
part in whole-stage code
+  generation, as filters, projections and aggregates do, Comet puts Spark's 
`ColumnarToRowExec`
+  between the two, and the generated code reads the cached Arrow vectors 
directly, with no
+  intermediate row. The plan shows this as a `ColumnarToRow` above the 
`InMemoryTableScan`. It
+  needs `spark.sql.inMemoryColumnarStorage.enableVectorizedReader` (on by 
default) and whole-stage
+  code generation, and applies to relations of at most 
`spark.sql.codegen.maxFields` fields (100 by
+  default, counting nested fields), beyond which Spark reads a cached relation 
only as rows. It is
+  not applied in plan-only mode (`spark.comet.explain.planOnly.enabled`), 
where Spark executes its

Review Comment:
   This sits in the bullet for native execution disabled, but plan-only mode 
only takes effect with `spark.comet.exec.enabled=true`. With execution off and 
`spark.comet.explain.planOnly.enabled=true`, the executed plan still gets the 
fused `ColumnarToRow`. Could this say that the fused reader is skipped in 
plan-only mode, which needs native execution on, rather than implying it 
applies here?



##########
docs/source/user-guide/latest/in-memory-cache.md:
##########
@@ -175,19 +176,45 @@ registrator.
 
 ## Limitations
 
-Reads that feed **Spark** operators rather than Comet ones are slower than 
Spark's own cache
-format, and the narrower the read, the wider the gap. Measured by the same 
benchmark over the same
-5M-row relation, with Comet off so that Spark operators consume the cached 
data:
-
-| Read shape              | Spark's cache format | Comet's cache format | 
Slowdown |
-| ----------------------- | -------------------: | -------------------: | 
-------: |
-| Row count only (0 of 6) |                35 ms |               183 ms |     
5.2x |
-| 1 of 6 columns          |                54 ms |               257 ms |     
4.8x |
-| 3 of 6 columns          |                98 ms |               331 ms |     
3.4x |
-| 6 of 6 columns          |               410 ms |               623 ms |     
1.5x |
-
-This is why the feature is off by default. The cause is not yet established;
-[#5485](https://github.com/apache/datafusion-comet/issues/5485) tracks it.
+Reads that feed **Spark** operators rather than Comet ones can be slower than 
Spark's own cache
+format, because every cached batch is decoded from Arrow before Spark reads 
it. How a Spark
+operator reads a relation cached in Comet's format depends on the scan below 
it:
+
+- With native execution enabled (`spark.comet.exec.enabled=true`), the cache 
is scanned by
+  `CometInMemoryTableScan`, and a Spark operator above it reads the scan's 
batches through
+  `CometColumnarToRow`, as it would above any other Comet operator.
+- With Comet enabled but native execution disabled, the cache is scanned by 
Spark's
+  `InMemoryTableScanExec`. When the operator directly above the scan takes 
part in whole-stage code
+  generation, as filters, projections and aggregates do, Comet puts Spark's 
`ColumnarToRowExec`
+  between the two, and the generated code reads the cached Arrow vectors 
directly, with no
+  intermediate row. The plan shows this as a `ColumnarToRow` above the 
`InMemoryTableScan`. It
+  needs `spark.sql.inMemoryColumnarStorage.enableVectorizedReader` (on by 
default) and whole-stage
+  code generation, and applies to relations of at most 
`spark.sql.codegen.maxFields` fields (100 by
+  default, counting nested fields), beyond which Spark reads a cached relation 
only as rows. It is
+  not applied in plan-only mode (`spark.comet.explain.planOnly.enabled`), 
where Spark executes its
+  own plan unchanged.
+- Otherwise the scan's row reader decodes each batch and writes its rows into 
one reused
+  `UnsafeRow`. That covers Comet or `spark.comet.exec.inMemoryCache.enabled` 
turned off at runtime,
+  and operators that do not take part in code generation, such as exchanges 
and limits, or a query
+  that returns the cached rows as they are.
+
+Measured by the same benchmark over the same 5M-row relation, with native 
execution off so that
+Spark operators consume the cached data, Comet disabled for the row reader and 
enabled for the fused
+reader (Apple M4, JDK 17, Spark 4.1; the average of two runs):
+
+| Read shape              | Spark's cache format | Comet's format, row reader 
| Comet's format, fused reader |
+| ----------------------- | -------------------: | -------------------------: 
| ---------------------------: |
+| Row count only (0 of 6) |                63 ms |                      57 ms 
|                        35 ms |
+| 1 of 6 columns          |                63 ms |                      77 ms 
|                        54 ms |
+| 3 of 6 columns          |               113 ms |                     176 ms 
|                       133 ms |
+| 6 of 6 columns          |               306 ms |                     500 ms 
|                       334 ms |
+
+The fused reader is faster than Spark's own format for the narrowest reads and 
within 20% of it for

Review Comment:
   The 20% holds for this benchmark's relation. On an all-bigint relation the 
fused reader took about 1.5 times as long as Spark's format even at six columns 
(summing every column: 80.8 against 53.2 ms). Could this sentence say it's for 
this relation, or mention that numeric-only relations see a larger gap?



##########
spark/src/main/scala/org/apache/comet/rules/CometCacheColumnarRule.scala:
##########
@@ -0,0 +1,100 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements.  See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership.  The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License.  You may obtain a copy of the License at
+ *
+ *   http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied.  See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+
+package org.apache.comet.rules
+
+import org.apache.spark.sql.catalyst.expressions.LeafExpression
+import org.apache.spark.sql.catalyst.expressions.codegen.CodegenFallback
+import org.apache.spark.sql.catalyst.rules.Rule
+import org.apache.spark.sql.comet.execution.arrow.ArrowCachedBatchSerializer
+import org.apache.spark.sql.execution.{CodegenSupport, ColumnarToRowExec, 
ColumnarToRowTransition, SparkPlan, WholeStageCodegenExec}
+import org.apache.spark.sql.execution.adaptive.QueryStageExec
+import org.apache.spark.sql.execution.columnar.InMemoryTableScanExec
+import org.apache.spark.sql.internal.SQLConf
+
+import org.apache.comet.CometConf.COMET_EXEC_IN_MEMORY_CACHE_ENABLED
+import org.apache.comet.CometSparkSessionExtensions.isCometLoaded
+
+/**
+ * Lets Spark's generated consumers read cached Arrow vectors without an 
intermediate UnsafeRow.
+ *
+ * Data flows upward. Spark's InputAdapter/whole-stage wrappers and an 
optional AQE cache stage
+ * are omitted:
+ * {{{
+ *   Before                              After
+ *   +------------------------+          +------------------------+
+ *   | Spark codegen consumer |          | Spark codegen consumer |
+ *   +------------------------+          +------------------------+
+ *               ^                                   ^
+ *               | UnsafeRow                         | column values
+ *   +------------------------+          +------------------------+
+ *   | InMemoryTableScanExec  |          | ColumnarToRowExec      |
+ *   | row iterator           |          | fused with consumer    |
+ *   +------------------------+          +------------------------+
+ *                                                   ^
+ *                                                   | ColumnarBatch
+ *                                       +------------------------+
+ *                                       | InMemoryTableScanExec  |
+ *                                       | Arrow vectors          |
+ *                                       +------------------------+
+ * }}}
+ *
+ * @param preview
+ *   true in the plan-only preview, which shows the plan Comet would execute. 
Otherwise the rule
+ *   leaves plans alone in plan-only mode, where Spark executes each query 
unchanged.
+ */
+case class CometCacheColumnarRule(preview: Boolean = false) extends 
Rule[SparkPlan] {
+  override def apply(plan: SparkPlan): SparkPlan = {
+    if (!isCometLoaded(conf) || !COMET_EXEC_IN_MEMORY_CACHE_ENABLED.get(conf)) 
return plan
+    if (!preview && CometRule.planOnlyApplies(conf, plan)) return plan
+    if (!conf.wholeStageEnabled) return plan
+    if (conf.getConf(SQLConf.CODEGEN_FACTORY_MODE).toString == "NO_CODEGEN") 
return plan
+
+    plan.transformUp {
+      case parent: CodegenSupport
+          if parent.supportCodegen && !parent.supportsColumnar &&
+            !parent.isInstanceOf[ColumnarToRowTransition] &&
+            !WholeStageCodegenExec.isTooManyFields(conf, parent.schema) &&

Review Comment:
   Near the top of the width this allows, the fused reader is slower than the 
row reader this PR adds. With Comet on and native execution off, summing every 
column of a cached all-bigint relation took 149.7 ms fused against 134.3 ms 
through `CachedBatchRowIterator` at 100 columns, and 141.2 against 130.5 ms at 
90 (24M values, interleaved, JDK 17, Spark 4.1). `count` over the same 100 
columns was 133.0 against 117.1 ms. At 50 and 75 columns the two were about 
even, and only at 6 columns did the fused reader clearly win (80.8 against 98.2 
ms). Main's row reader is within a couple of percent of the new one at these 
widths, so here the rule makes the read about 10% slower than today. It isn't 
the huge-method limit, since the fused stage's largest method was 6,155 bytes 
at 100 columns. Could the rule stop fusing at a narrower scan output than 
`spark.sql.codegen.maxFields`, or could the wide benchmark get a fused arm at 
100 columns so the crossover is visible?



##########
docs/source/user-guide/latest/in-memory-cache.md:
##########
@@ -175,19 +176,45 @@ registrator.
 
 ## Limitations
 
-Reads that feed **Spark** operators rather than Comet ones are slower than 
Spark's own cache
-format, and the narrower the read, the wider the gap. Measured by the same 
benchmark over the same
-5M-row relation, with Comet off so that Spark operators consume the cached 
data:
-
-| Read shape              | Spark's cache format | Comet's cache format | 
Slowdown |
-| ----------------------- | -------------------: | -------------------: | 
-------: |
-| Row count only (0 of 6) |                35 ms |               183 ms |     
5.2x |
-| 1 of 6 columns          |                54 ms |               257 ms |     
4.8x |
-| 3 of 6 columns          |                98 ms |               331 ms |     
3.4x |
-| 6 of 6 columns          |               410 ms |               623 ms |     
1.5x |
-
-This is why the feature is off by default. The cause is not yet established;
-[#5485](https://github.com/apache/datafusion-comet/issues/5485) tracks it.
+Reads that feed **Spark** operators rather than Comet ones can be slower than 
Spark's own cache
+format, because every cached batch is decoded from Arrow before Spark reads 
it. How a Spark
+operator reads a relation cached in Comet's format depends on the scan below 
it:
+
+- With native execution enabled (`spark.comet.exec.enabled=true`), the cache 
is scanned by
+  `CometInMemoryTableScan`, and a Spark operator above it reads the scan's 
batches through
+  `CometColumnarToRow`, as it would above any other Comet operator.
+- With Comet enabled but native execution disabled, the cache is scanned by 
Spark's
+  `InMemoryTableScanExec`. When the operator directly above the scan takes 
part in whole-stage code
+  generation, as filters, projections and aggregates do, Comet puts Spark's 
`ColumnarToRowExec`
+  between the two, and the generated code reads the cached Arrow vectors 
directly, with no
+  intermediate row. The plan shows this as a `ColumnarToRow` above the 
`InMemoryTableScan`. It
+  needs `spark.sql.inMemoryColumnarStorage.enableVectorizedReader` (on by 
default) and whole-stage
+  code generation, and applies to relations of at most 
`spark.sql.codegen.maxFields` fields (100 by
+  default, counting nested fields), beyond which Spark reads a cached relation 
only as rows. It is
+  not applied in plan-only mode (`spark.comet.explain.planOnly.enabled`), 
where Spark executes its
+  own plan unchanged.
+- Otherwise the scan's row reader decodes each batch and writes its rows into 
one reused
+  `UnsafeRow`. That covers Comet or `spark.comet.exec.inMemoryCache.enabled` 
turned off at runtime,
+  and operators that do not take part in code generation, such as exchanges 
and limits, or a query

Review Comment:
   `CollectLimitExec` doesn't take part in code generation, but 
`LocalLimitExec` and `GlobalLimitExec` do, so a limit inside a larger plan 
reads through the fused reader. Should this say a top-level limit?



-- 
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