Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
29 commits
Select commit Hold shift + click to select a range
d7ae3bb
perf: project cached batches by buffer selection, prune on collated s…
andygrove Aug 28, 2026
f59c9dc
refactor: use Spark's interpreted ordering for bounds, hoist projecti…
andygrove Aug 29, 2026
ccd469e
fix: relocate the arrow-compression service file when shading
andygrove Aug 29, 2026
e71d802
test: drop the cache leak test that depends on zstd corruption detection
andygrove Aug 29, 2026
d8d4196
feat: enable Comet's in-memory cache by default
andygrove Sep 2, 2026
4010f78
test: cover nested columns in the cached-batch projection tests and b…
andygrove Sep 5, 2026
8c19267
fix: drop a redundant string interpolator flagged by scalafix Redunda…
andygrove Sep 5, 2026
1699665
Merge remote-tracking branch 'apache/main' into feat/cache-buffer-sel…
andygrove Sep 7, 2026
792a465
Merge branch 'main' into feat/cache-enabled-by-default
andygrove Sep 8, 2026
bac454e
Merge remote-tracking branch 'apache/main' into feat/cache-buffer-sel…
cincrement Sep 10, 2026
f05c204
Merge branch 'main' into feat/cache-buffer-selection-projection
andygrove Sep 14, 2026
3e74e9d
review: check the cached layout, and address the rest of the review
andygrove Sep 15, 2026
b978e54
review: own the write-side compression buffers, fix the activation ex…
andygrove Sep 15, 2026
dbf487b
Merge remote-tracking branch 'origin/feat/cache-buffer-selection-proj…
andygrove Sep 21, 2026
a71e8cb
Merge remote-tracking branch 'apache/main' into feat/cache-enabled-by…
andygrove Sep 21, 2026
e483d08
fix: write cached batches to the schema width, not the batch width
andygrove Sep 21, 2026
3cf15ac
Merge branch 'fix/cache-wide-columnar-batch' into feat/cache-enabled-…
andygrove Sep 21, 2026
c7f1ce3
Merge branch 'main' into feat/cache-enabled-by-default
andygrove Sep 22, 2026
89cd107
Merge remote-tracking branch 'apache/main' into HEAD
andygrove Sep 24, 2026
d5983c2
Merge remote-tracking branch 'apache/main' into feat/cache-enabled-by…
andygrove Sep 24, 2026
299d381
Merge remote-tracking branch 'apache/main' into feat/cache-enabled-by…
andygrove Sep 25, 2026
a371d64
fix: report a cached relation's decoded size to the planner, not its …
andygrove Sep 29, 2026
6899380
fix: keep Spark's cache scan for a relation whose cached plan records…
andygrove Sep 29, 2026
d84b1f0
Merge remote-tracking branch 'apache/main' into feat/cache-enabled-by…
andygrove Sep 30, 2026
0ba6716
test: install Comet's cache serializer in the Spark SQL test sessions
andygrove Sep 30, 2026
37560b8
test: fix three Spark SQL test adaptations for Comet's cache format
andygrove Sep 30, 2026
c57eb9d
Merge branch 'fix/cache-decoded-size-stats' into feat/cache-enabled-b…
andygrove Oct 1, 2026
75b714b
Merge branch 'fix/cache-observed-metrics' into feat/cache-enabled-by-…
andygrove Oct 1, 2026
45dcd0e
Merge branch 'main' into feat/cache-enabled-by-default
andygrove Oct 1, 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
327 changes: 316 additions & 11 deletions dev/diffs/3.4.3.diff

Large diffs are not rendered by default.

370 changes: 356 additions & 14 deletions dev/diffs/3.5.9.diff

Large diffs are not rendered by default.

416 changes: 400 additions & 16 deletions dev/diffs/4.0.4.diff

Large diffs are not rendered by default.

489 changes: 466 additions & 23 deletions dev/diffs/4.1.3.diff

Large diffs are not rendered by default.

489 changes: 466 additions & 23 deletions dev/diffs/4.2.0.diff

Large diffs are not rendered by default.

19 changes: 13 additions & 6 deletions docs/source/user-guide/latest/in-memory-cache.md
Original file line number Diff line number Diff line change
Expand Up @@ -24,13 +24,13 @@ format that Comet operators read directly. Without it, a cached table is stored
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.

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**. To turn it off, set the config at startup,
alongside the rest of Comet's configuration:
Comment on lines +27 to +28

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This key did not exist in 1.0.0, so a user upgrading from 1.0.0 goes from Spark's cache format to Comet's without setting anything. With spark.kryo.registrationRequired=true and no CometKryoRegistrator, a df.cache() that spills to disk now fails with "Class is not registered" where it did not before. The plugin only logs a warning for that.

The versioning policy counts a new error under the same explicit configuration as a behavior change. Could you add an entry to the upgrade guide under the next release that covers the format change and the Kryo requirement? The policy asks for a spark.comet.legacy.* key, but spark.comet.exec.inMemoryCache.enabled=false already restores the old behavior, so naming that key in the entry seems enough. If you read the policy differently, it would be good to settle that here, since this is one of the first behavior changes since 1.0.0.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Agreed on the upgrade guide entry, covering both the format change and the Kryo requirement. Moving the flip past 1.1.0 changes one premise, though. 1.1.0 ships this key with a default of false, so turning it on in the next release is a change to an existing key's default, which is the first case the policy lists. Let's settle the legacy-key question when this comes out of draft.


```shell
$SPARK_HOME/bin/spark-shell \
... \
--conf spark.comet.exec.inMemoryCache.enabled=true
--conf spark.comet.exec.inMemoryCache.enabled=false
```

It has to be set before the `SparkContext` starts. Comet's driver plugin chooses
Expand All @@ -48,6 +48,9 @@ With Comet's serializer installed as `spark.sql.cache.serializer`:
- Cached tables are scanned by `CometInMemoryTableScan`, which feeds Comet operators directly.
- Per-batch column statistics are recorded in the layout Spark's `SimpleMetricsCachedBatchSerializer`
expects, so Spark can prune whole cached batches on a predicate before any of them is decoded.
- The size Spark's planner sees for a cached relation is its decoded Arrow size, not the compressed
size it occupies in memory, as with Spark's own cache formats. Compression therefore does not
change how queries over a cached relation are planned, such as whether a join broadcasts it.

Relations whose schema Comet's Arrow writer cannot store — interval types, most notably — are
delegated in full to Spark's default cache format, per relation. Which format a relation uses does
Expand All @@ -58,6 +61,10 @@ 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.

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
metrics only through that scan.

## Storage format

Each cached batch is stored as a single Arrow IPC record batch message and its body.
Expand Down Expand Up @@ -102,7 +109,7 @@ nowhere to record either that a column is dictionary encoded or the dictionary i

| Config | Default | Description |
| ------------------------------------------------------- | ------- | ---------------------------------------------------------------------------------------------------------------------------------------------- |
| `spark.comet.exec.inMemoryCache.enabled` | `false` | Whether to store and scan Spark's in-memory cache in Comet's format. Read at startup. |
| `spark.comet.exec.inMemoryCache.enabled` | `true` | Whether to store and scan Spark's in-memory cache in Comet's format. Read at startup. |
| `spark.comet.exec.inMemoryCache.compression.codec` | `zstd` | Arrow IPC compression codec for cached data: `zstd` or `none`. Affects newly cached data only — a batch records the codec it was written with. |
| `spark.comet.exec.inMemoryCache.compression.zstd.level` | `1` | Compression level when the codec is `zstd`. Ignored otherwise. |

Expand Down Expand Up @@ -190,8 +197,8 @@ format, and the narrower the read, the wider the gap. Measured by the same bench
| 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.
This is the main reason the feature is still described as experimental. The cause is not yet
established; [#5485](https://github.com/apache/datafusion-comet/issues/5485) tracks it.

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
2 changes: 1 addition & 1 deletion spark/src/main/scala/org/apache/comet/CometConf.scala
Original file line number Diff line number Diff line change
Expand Up @@ -280,7 +280,7 @@ object CometConf extends ShimCometConf {
"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)
.createWithDefault(true)

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

#5485 treats the slower Spark-operator reads as acceptable because "the cache path is off by default ... rather than a regression in a shipped path." This PR makes it a shipped path, and the format is fixed when the relation materializes, so a user can't avoid it for one query.

With the plugin gated as suggested above, the remaining exposure is a query where the cached scan runs natively and a Spark operator above it reads through a columnar-to-row transition, and a session that turns Comet off at runtime after caching. Is there a benchmark number for the first case against Spark's own cache format? The published numbers compare against Comet off entirely. For the second case, option 1 in #5485 (a fallback reason when Spark operators read a relation stored in Comet's format) is what would tell a user to turn the feature off. Could that land before or with this PR?

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

The audit turned up part of the answer. Under AQE the first case is the normal outcome, not an edge case. Once the table-cache stage materializes, the re-plan leaves the operators above it on Spark, so a plain aggregate or join over a cached table reads Comet's format through a CometColumnarToRow (#6202). CometInMemoryCacheBenchmark runs with AQE off, so the published numbers don't show it. I'll fix #6202 first and then benchmark with AQE on against Spark's own format, so the number measures the path users will actually get.

Yes to option 1 from #5485 landing with the default flip. It will need to cover the #6202 path too, since nothing records a fallback reason there today.

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

#6202 is fixed by #6208, which is now merged into this branch in 299d381. Under AQE, an aggregate or join over a cached table now stays native once the table-cache stage materializes, so the first case is no longer the normal outcome: it takes an operator Comet does not support above the cached scan. Option 1 no longer has a #6202 path to cover either, since the operators above the stage now convert, or record their own fallback reason, like any other operator. The benchmark with AQE on against Spark's own format is next.


val COMET_EXEC_IN_MEMORY_CACHE_COMPRESSION_CODEC: ConfigEntry[String] =
conf("spark.comet.exec.inMemoryCache.compression.codec")
Expand Down
14 changes: 12 additions & 2 deletions spark/src/main/scala/org/apache/comet/rules/CometExecRule.scala
Original file line number Diff line number Diff line change
Expand Up @@ -386,8 +386,12 @@ case class CometExecRule(session: SparkSession)
val cometCacheFormat = usesCometCacheSerializer &&
ArrowCachedBatchSerializer.supportsSchema(scan.relation.output)
val nativeCacheEnabled = CometConf.COMET_EXEC_IN_MEMORY_CACHE_ENABLED.get(conf)
// Walks the cached plan, so it is only consulted once the native scan is otherwise
// possible. See CometInMemoryTableScanExec.recordsObservedMetrics.
val nativeScan = nativeCacheEnabled && cometCacheFormat &&
!CometInMemoryTableScanExec.recordsObservedMetrics(scan.relation)

if (nativeCacheEnabled && cometCacheFormat) {
if (nativeScan) {
convertToComet(scan, CometInMemoryTableScanExec).getOrElse(scan)
} else {
// The native cache scan is not available for this relation. Record why, then take the
Expand All @@ -398,14 +402,20 @@ case class CometExecRule(session: SparkSession)
scan,
s"Comet in-memory cache requires ${classOf[ArrowCachedBatchSerializer].getName} " +
s"but this relation was cached with ${serializer.getClass.getName}")
} else if (nativeCacheEnabled) {
} else if (nativeCacheEnabled && !cometCacheFormat) {
val unsupported = scan.relation.output
.filterNot(a => ArrowCachedBatchSerializer.supportsType(a.dataType))
.map(a => s"${a.name}: ${a.dataType.simpleString}")
withFallbackReason(
scan,
"Comet in-memory cache does not support the type of these cached columns, so the " +
s"relation was cached in Spark's default format: ${unsupported.mkString(", ")}")
} else if (nativeCacheEnabled) {
withFallbackReason(
scan,
"Comet in-memory cache does not scan a relation whose cached plan records " +
"Dataset.observe metrics, because Spark collects those metrics only through " +
"InMemoryTableScanExec")
} else if (usesCometCacheSerializer) {
withFallbackReason(
scan,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -26,8 +26,9 @@ import org.apache.spark.sql.catalyst.expressions.Attribute
import org.apache.spark.sql.catalyst.plans.logical.Statistics
import org.apache.spark.sql.columnar.{CachedBatch, CachedBatchSerializer}
import org.apache.spark.sql.comet.shims.ShimCometInMemoryTableScanExec
import org.apache.spark.sql.execution.SparkPlan
import org.apache.spark.sql.execution.columnar.{CachedRDDBuilder, InMemoryTableScanExec}
import org.apache.spark.sql.execution.{CollectMetricsExec, SparkPlan}
import org.apache.spark.sql.execution.adaptive.AdaptiveSparkPlanHelper
import org.apache.spark.sql.execution.columnar.{CachedRDDBuilder, InMemoryRelation, InMemoryTableScanExec}
import org.apache.spark.sql.execution.metric.{SQLMetric, SQLMetrics}
import org.apache.spark.sql.vectorized.ColumnarBatch

Expand Down Expand Up @@ -170,4 +171,24 @@ object CometInMemoryTableScanExec extends CometOperatorSerde[InMemoryTableScanEx
op.output))
}

/**
* Whether `relation`'s cached plan records observed metrics, from `Dataset.observe`.
*
* Spark collects those metrics once a query finishes, with `CollectMetricsExec.collect`, and
* that reaches the ones recorded inside a cached plan only through an `InMemoryTableScanExec`
* over it. It does not know this node, so replacing the scan of such a relation leaves its
* metrics empty, and on Spark 3.4 leaves `Observation.get` waiting for good. The walk mirrors
* `CollectMetricsExec.collect`, through subqueries, adaptive plans and nested caches.
*/
def recordsObservedMetrics(relation: InMemoryRelation): Boolean =
ObservedMetrics.recordedIn(relation.cachedPlan)

private object ObservedMetrics extends AdaptiveSparkPlanHelper {
def recordedIn(plan: SparkPlan): Boolean =
collectWithSubqueries(plan) {
case _: CollectMetricsExec => true
case scan: InMemoryTableScanExec => recordedIn(scan.relation.cachedPlan)
}.contains(true)
}

}
Original file line number Diff line number Diff line change
Expand Up @@ -52,10 +52,14 @@ import org.apache.comet.vector.NativeUtil
* and length, so `CachedBatchIpc.Projection.load` copies out just the selected columns' byte
* ranges. The cache manager still owns storage and eviction; this class only changes the cached
* payload.
*
* `sizeInBytes` is not the payload's size. It is inherited from `SimpleMetricsCachedBatch`, which
* sums the per-column sizes in `stats`, and those are decoded sizes. Spark's planner reads it as
* the size of a materialized cached relation, as it does for Spark's own cache formats, which
* also report decoded sizes. See `statsRow`.
*/
private case class CometCachedBatch(
override val numRows: Int,
override val sizeInBytes: Long,
override val stats: InternalRow,
bytes: ChunkedByteBuffer)
extends SimpleMetricsCachedBatch
Expand Down Expand Up @@ -350,11 +354,13 @@ class ArrowCachedBatchSerializer extends SimpleMetricsCachedBatchSerializer {
values(base + 1) = upper(c)
values(base + 2) = nulls(c)
values(base + 3) = numRows
// The stored size of the column's own Arrow buffers, taken from the message's buffer
// layout, so it is exact rather than an estimate. Cache pruning uses
// bounds/null-count/row-count rather than this field, but Spark reserves it and reports it,
// so record the real value. The per-batch message framing is not attributed to any column,
// so these sum to slightly less than sizeInBytes.
// The column's decoded size: the plain length of its own Arrow buffers before compression,
// which is what Spark's own Arrow cache format records here. SimpleMetricsCachedBatch sums
// these into sizeInBytes, which Spark's planner reads as the size of a materialized cached
// relation, for the broadcast threshold and the shuffled hash join build side among others.
// The compressed payload can be several times smaller, and a relation reported at that size
// would be planned differently from the same relation in Spark's cache format, for example
// broadcast where Spark's cache would have it shuffled.
values(base + 4) = columnSizes(c)
c += 1
}
Expand Down Expand Up @@ -423,7 +429,6 @@ class ArrowCachedBatchSerializer extends SimpleMetricsCachedBatchSerializer {

CometCachedBatch(
numRows = numRows,
sizeInBytes = bytes.size,
stats = statsRow(lower, upper, nulls, numRows, columnSizes),
bytes = bytes)
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -186,9 +186,10 @@ private[comet] object CachedBatchIpc {
* batch arrives at whatever size the plan above produced. Chunks are appended rather than grown
* and recopied, so the write also never holds the payload twice.
*
* Returns the message and the on-body compressed size of each top-level column, which the
* caller records in the statistics row. The sizes come from the message's own buffer layout, so
* they are the real stored sizes rather than an estimate.
* Returns the message and the decoded size of each top-level column, which the caller records
* in the statistics row. Each size is measured on the batch before compression, from the plain
* lengths of the column's own buffers, which is what `getBufferSize` reports for a vector and
* what Spark's own Arrow cache format records.
*
* Dictionary-encoded columns are decoded to their plain form first. A payload with no Schema
* message cannot describe a dictionary encoding, and the schema the reader rebuilds from Spark
Expand All @@ -215,16 +216,18 @@ private[comet] object CachedBatchIpc {

// Unloaded plain and compressed afterwards rather than by handing the codec to the unloader;
// see compressed for why.
val fields = vectors.map(_.getField)
val unloader = new VectorUnloader(root, true, NoCompressionCodec.INSTANCE, true)
val plainBatch = unloader.getRecordBatch
val recordBatch =
try compressed(plainBatch, codec, allocator)
finally plainBatch.close()
val (sizes, recordBatch) =
try {
(columnSizes(fields, plainBatch), compressed(plainBatch, codec, allocator))
} finally {
plainBatch.close()
}
try {
val fields = vectors.map(_.getField)
// Leaves the batch in the state serializeBatches leaves one. The record batch holds its
// own buffers by now, so this does not touch it, and getField still answers afterwards:
// clearing releases buffers, not the schema.
// own buffers by now, so this does not touch it.
root.clear()

val out = new ChunkedByteBufferOutputStream(chunkSize, ByteBuffer.allocate)
Expand All @@ -234,7 +237,7 @@ private[comet] object CachedBatchIpc {
} finally {
out.close()
}
(out.toChunkedByteBuffer, columnSizes(fields, recordBatch))
(out.toChunkedByteBuffer, sizes)
} finally {
recordBatch.close()
}
Expand Down Expand Up @@ -608,11 +611,10 @@ private[comet] object CachedBatchIpc {
}

/**
* The on-body compressed size of each top-level column.
* The size of each top-level column's buffers in `recordBatch`.
*
* Each column owns the run of buffers its subtree occupies, so its stored size is the sum of
* those buffers' recorded lengths. With one payload per batch these are the only per-column
* sizes available -- there is no separate stream to measure -- and they are exact.
* Each column owns the run of buffers its subtree occupies, so its size is the sum of those
* buffers' recorded lengths.
*/
private def columnSizes(fields: Seq[Field], recordBatch: ArrowRecordBatch): Array[Long] = {
val buffers = recordBatch.getBuffersLayout
Expand Down
Loading
Loading