Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
23 commits
Select commit Hold shift + click to select a range
ec92ee5
perf: reduce Spark cache row conversion overhead
peterxcli Sep 11, 2026
17dcdc6
chore: remove cache row reader benchmark results
peterxcli Sep 11, 2026
8dc61ad
perf: feed cached Arrow columns into Spark codegen
peterxcli Sep 11, 2026
091eb00
docs: illustrate Comet cache columnar rewrite
peterxcli Sep 11, 2026
eafdba5
fix: honor Comet disable switches for fused cache reads
peterxcli Sep 13, 2026
49b7ec0
fix: align cache fusion with Spark codegen settings
peterxcli Sep 15, 2026
03981e5
Merge branch 'main' into codex/cache-spark-consumer-benchmark
peterxcli Sep 15, 2026
ee0e24f
ci: retry after runner and dependency download failures
peterxcli Sep 15, 2026
14b5d7b
Merge branch 'main' into codex/cache-spark-consumer-benchmark
peterxcli Sep 16, 2026
cea2caa
Merge upstream main into codex/cache-spark-consumer-benchmark
peterxcli Sep 21, 2026
eadfdff
Merge upstream main into codex/cache-spark-consumer-benchmark
peterxcli Sep 22, 2026
64a25bc
Merge remote-tracking branch 'upstream/main' into HEAD
peterxcli Sep 25, 2026
dbfb182
Merge remote-tracking branch 'upstream/main' into HEAD
peterxcli Sep 25, 2026
3aa6030
perf: split the generated cache row reader for wide projections
peterxcli Sep 28, 2026
351503a
fix: do not fuse cache reads in plan-only mode
peterxcli Sep 28, 2026
d4ad3c4
test: check fused cache reads with Comet on and native execution off
peterxcli Sep 28, 2026
f6f245f
test: measure the fused cache reader in CometInMemoryCacheBenchmark
peterxcli Sep 28, 2026
384ea00
docs: describe how Spark operators read Comet's cache format
peterxcli Sep 28, 2026
a150075
Merge remote-tracking branch 'apache/main' into HEAD
andygrove Sep 29, 2026
4689053
Merge upstream main into codex/cache-spark-consumer-benchmark
peterxcli Sep 30, 2026
869c5d2
fix: avoid numeric widening in cache benchmark
peterxcli Sep 30, 2026
0da3690
ci: retry preflight after Maven Central download failure
peterxcli Sep 30, 2026
73092ba
Merge remote-tracking branch 'upstream/main' into HEAD
peterxcli Oct 2, 2026
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
1 change: 1 addition & 0 deletions .github/workflows/pr_build_linux.yml
Original file line number Diff line number Diff line change
Expand Up @@ -571,6 +571,7 @@ jobs:
org.apache.spark.sql.CometSortCollationSuite
org.apache.comet.CometFuzzAggregateSuite
org.apache.spark.sql.comet.execution.arrow.CometArrowStreamSuite
org.apache.spark.sql.comet.execution.arrow.CachedBatchRowIteratorSuite
org.apache.spark.sql.comet.execution.arrow.CometStringWriterSuite
org.apache.spark.sql.CometSparkInternalFunctionsSuite
- name: "expressions"
Expand Down
1 change: 1 addition & 0 deletions .github/workflows/pr_build_macos.yml
Original file line number Diff line number Diff line change
Expand Up @@ -275,6 +275,7 @@ jobs:
org.apache.spark.sql.CometSortCollationSuite
org.apache.comet.CometFuzzAggregateSuite
org.apache.spark.sql.comet.execution.arrow.CometArrowStreamSuite
org.apache.spark.sql.comet.execution.arrow.CachedBatchRowIteratorSuite
org.apache.spark.sql.comet.execution.arrow.CometStringWriterSuite
org.apache.spark.sql.CometSparkInternalFunctionsSuite
- name: "expressions"
Expand Down
55 changes: 41 additions & 14 deletions docs/source/user-guide/latest/in-memory-cache.md
Original file line number Diff line number Diff line change
Expand Up @@ -59,7 +59,8 @@ relation whose format could change mid-session could not be read back reliably.
codec is a runtime config, but each batch records the codec it was written with, so data cached
under one setting stays readable after the setting changes. Turning
`spark.comet.exec.inMemoryCache.enabled` off at runtime only sends cached scans back to Spark's
execution path; the cached data stays readable either way.
execution path, where Spark operators read them through the row reader described under
[Limitations](#limitations); the cached data stays readable either way.

A relation whose cached plan records observed metrics, from `Dataset.observe`, is still stored in
Comet's format but is scanned by Spark's `InMemoryTableScanExec`, because Spark collects those
Expand Down Expand Up @@ -186,19 +187,45 @@ registrator. Native broadcast needs the same registrator even when the cache is

## 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
Comment thread
peterxcli marked this conversation as resolved.
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
Comment thread
peterxcli marked this conversation as resolved.
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
Comment thread
peterxcli marked this conversation as resolved.
the others. The row reader takes up to 1.6 times as long, and it is the only reader for relations
wider than `spark.sql.codegen.maxFields`: reading every column of relations of 100, 200 and 1500
nullable `bigint` columns took 2.2 to 2.5 times as long as from Spark's format. These gaps are why
the feature is still off by default;
[#5485](https://github.com/apache/datafusion-comet/issues/5485) tracks them.

Comet's serializer exists because Spark's own Arrow cache format
([SPARK-57268](https://issues.apache.org/jira/browse/SPARK-57268)) is only available from Spark
Expand Down
37 changes: 21 additions & 16 deletions spark/src/main/scala/org/apache/comet/CometConf.scala
Original file line number Diff line number Diff line change
Expand Up @@ -263,22 +263,27 @@ object CometConf extends ShimCometConf {
val COMET_EXEC_IN_MEMORY_CACHE_ENABLED: ConfigEntry[Boolean] =
conf("spark.comet.exec.inMemoryCache.enabled")
.category(CATEGORY_EXEC)
.doc("Whether to enable Comet native execution for in-memory cached tables. Its value at " +
"startup also decides whether CometDriverPlugin installs Comet's cache serializer, " +
"which stores cached data in Arrow format. The plugin installs it only if " +
"spark.comet.enabled and spark.comet.exec.enabled are also enabled at startup. " +
"Because spark.sql.cache.serializer is a " +
"static config, the cached format is fixed for the application, and disabling this " +
"at runtime only sends cached scans back to Spark's execution path. Relations whose " +
"schema Comet's Arrow writer does not support are always cached in Spark's default " +
"format. Each cached batch is stored as one Arrow IPC record batch with per-buffer " +
"zstd compression, and a scan copies out only the buffers of the columns it projected, " +
"so the unselected ones are never decompressed. Reads that feed Spark operators rather " +
"than Comet ones still pay a row conversion the default format avoids, and can be " +
"slower than Spark's cache. With spark.kryo.registrationRequired=true, also set " +
"spark.kryo.registrator=org.apache.comet.CometKryoRegistrator before creating the " +
"SparkContext, otherwise caching fails as soon as a block is serialized, including " +
"the disk half of the default MEMORY_AND_DISK storage level.")
.doc(
"Whether to enable Comet native scans and fused Spark reads of in-memory cached tables. " +
"Requires spark.comet.enabled=true. At startup, this setting also decides whether " +
"CometDriverPlugin installs Comet's cache serializer, which stores cached data in " +
"Arrow format. The plugin installs it only if spark.comet.enabled and " +
"spark.comet.exec.enabled are also enabled at startup. " +
"Because spark.sql.cache.serializer is a " +
"static config, the cached format is fixed for the application, and disabling this " +
"or spark.comet.enabled at runtime sends cached scans back to Spark's execution path " +
"without the fused reader. Relations whose schema Comet's Arrow writer does not " +
"support are always cached in Spark's default " +
"format. Each cached batch is stored as one Arrow IPC record batch with per-buffer " +
"zstd compression, and a scan copies out only the buffers of the columns it " +
"projected, so the unselected ones are never decompressed. Eligible Spark " +
"whole-stage codegen consumers read cached vectors directly when vectorized cache " +
"reading is enabled; other Spark row consumers use a reusable row buffer. Decoding " +
"costs can still make wide numeric reads slower than Spark's default cache. With " +
Comment thread
peterxcli marked this conversation as resolved.
"spark.kryo.registrationRequired=true, also set " +
"spark.kryo.registrator=org.apache.comet.CometKryoRegistrator before creating the " +
"SparkContext, otherwise caching fails as soon as a block is serialized, including " +
"the disk half of the default MEMORY_AND_DISK storage level.")
.booleanConf
.createWithDefault(false)

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -53,7 +53,7 @@ import org.apache.comet.shims.ShimCometSparkSessionExtensions
* CometSubqueryBroadcastExec for exchange reuse with Comet broadcasts
* b. insertTransitions: ColumnarToRow/RowToColumnar added
* c. postColumnarTransitions: RevertNativeForTransitionHeavyStages,
* EliminateRedundantTransitions
* EliminateRedundantTransitions, CometCacheColumnarRule
* 5. ReuseExchangeAndSubquery -- Spark deduplicates subqueries (sees Comet nodes)
* }}}
*
Expand All @@ -77,7 +77,7 @@ import org.apache.comet.shims.ShimCometSparkSessionExtensions
* a. preColumnarTransitions: CometRule (no-op, already converted)
* b. insertTransitions
* c. postColumnarTransitions: RevertNativeForTransitionHeavyStages,
* EliminateRedundantTransitions
* EliminateRedundantTransitions, CometCacheColumnarRule
* }}}
*
* On Spark 3.4, injectQueryStageOptimizerRule is unavailable. CometExecRule does not wrap SABs,
Expand Down
Original file line number Diff line number Diff line change
@@ -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
Comment thread
peterxcli marked this conversation as resolved.
if (!preview && CometRule.planOnlyApplies(conf, plan)) return plan
if (!conf.wholeStageEnabled) return plan
Comment thread
peterxcli marked this conversation as resolved.
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) &&
Comment thread
peterxcli marked this conversation as resolved.
!parent.children.exists(p => WholeStageCodegenExec.isTooManyFields(conf, p.schema)) &&
!parent.expressions.exists(_.exists {
case _: LeafExpression => false
case _: CodegenFallback => true
case _ => false
}) =>
// Match the consuming edge rather than every scan: an existing columnar consumer (or a
// cache stage being materialized by AQE) must keep receiving batches. Spark inserts an
// InputAdapter around the scan later, while this transition fuses with the row consumer.
parent.withNewChildren(parent.children.map {
case child if isColumnarCometCache(child) => ColumnarToRowExec(child)
case child => child
})
}
}

private def isColumnarCometCache(plan: SparkPlan): Boolean = {
plan.supportsColumnar && (plan match {
case scan: InMemoryTableScanExec =>
// The serializer delegates unsupported schemas to Spark, whose cache keeps its own reader.
scan.relation.cacheBuilder.serializer.isInstanceOf[ArrowCachedBatchSerializer] &&
ArrowCachedBatchSerializer.supportsSchema(scan.relation.output)
case stage: QueryStageExec => isColumnarCometCache(stage.plan)
case _ => false
})
}
}
38 changes: 25 additions & 13 deletions spark/src/main/scala/org/apache/comet/rules/CometRule.scala
Original file line number Diff line number Diff line change
Expand Up @@ -29,18 +29,37 @@ import org.apache.spark.sql.execution.{ApplyColumnarRulesAndInsertTransitions, B
import org.apache.spark.sql.execution.adaptive.{AdaptiveSparkPlanExec, InsertAdaptiveSparkPlan, QueryStageExec}
import org.apache.spark.sql.execution.exchange.{BroadcastExchangeExec, Exchange}
import org.apache.spark.sql.execution.reuse.ReuseExchangeAndSubquery
import org.apache.spark.sql.internal.SQLConf

import org.apache.comet.{CometConf, ExtendedExplainInfo}
import org.apache.comet.CometSparkSessionExtensions.isCometLoaded
import org.apache.comet.shims.ShimCometStreaming

object CometRule {

/** Comet's post-columnar rules, shared by `CometColumnar` and the plan-only preview. */
def postColumnarRules(session: SparkSession, wholePlan: Boolean = false): Seq[Rule[SparkPlan]] =
/**
* Comet's post-columnar rules, shared by `CometColumnar` and the plan-only preview.
*
* @param preview
* true for the plan-only preview, which holds the whole plan and shows the plan Comet would
* execute.
*/
def postColumnarRules(session: SparkSession, preview: Boolean = false): Seq[Rule[SparkPlan]] =
Seq(
RevertNativeForTransitionHeavyStages(session, wholePlan),
EliminateRedundantTransitions(session))
RevertNativeForTransitionHeavyStages(session, wholePlan = preview),
EliminateRedundantTransitions(session),
CometCacheColumnarRule(preview))

/**
* Whether plan-only mode applies to `plan`, so that Comet only reports the plan it would
* execute and Spark executes `plan` unchanged. Mirrors the conversion rules' own guards;
* plan-only is scoped to exec being enabled.
*/
private[comet] def planOnlyApplies(conf: SQLConf, plan: SparkPlan): Boolean =
CometConf.COMET_EXPLAIN_PLAN_ONLY_ENABLED.get(conf) &&
isCometLoaded(conf) &&
!ShimCometStreaming.isStreamingPlan(plan) &&
CometConf.COMET_EXEC_ENABLED.get(conf)

/**
* Canonical hashes of the subquery plans reported for the query this thread is preparing. Spark
Expand Down Expand Up @@ -140,7 +159,7 @@ case class CometRule(session: SparkSession, queryStagePrep: Boolean = false)
private val execRule = CometExecRule(session)

override def apply(plan: SparkPlan): SparkPlan = {
if (planOnlyApplies(plan)) {
if (CometRule.planOnlyApplies(conf, plan)) {
reportPlanOnlyCoverage(plan)
plan
} else {
Expand All @@ -150,13 +169,6 @@ case class CometRule(session: SparkSession, queryStagePrep: Boolean = false)

private def convert(plan: SparkPlan): SparkPlan = execRule.apply(scanRule.apply(plan))

/** Mirrors the conversion rules' own guards; plan-only is scoped to exec being enabled. */
private def planOnlyApplies(plan: SparkPlan): Boolean =
CometConf.COMET_EXPLAIN_PLAN_ONLY_ENABLED.get(conf) &&
isCometLoaded(conf) &&
!ShimCometStreaming.isStreamingPlan(plan) &&
CometConf.COMET_EXEC_ENABLED.get(conf)

/** Logs the Comet plan for `plan` unless already reported. Never fails the query. */
private def reportPlanOnlyCoverage(plan: SparkPlan): Unit = {
try {
Expand All @@ -183,7 +195,7 @@ case class CometRule(session: SparkSession, queryStagePrep: Boolean = false)
val withTransitions =
ApplyColumnarRulesAndInsertTransitions(Seq.empty, outputsColumnar = false).apply(converted)
val preview = CometRule
.postColumnarRules(session, wholePlan = true)
.postColumnarRules(session, preview = true)
.foldLeft(withTransitions) { case (p, rule) => rule(p) }
if (topLevel) ReuseExchangeAndSubquery.apply(preview) else preview
}
Expand Down
Loading
Loading