Memory Tuning#
Comet needs two things configured: an off-heap pool for it to draw its reservations from, and enough executor memory overhead to cover the part of its footprint that no pool tracks. See Configuring Comet Memory and Configuring Executor Memory Overhead. To find out how much memory Comet uses on a real workload, see Sizing the Overhead from the Memory Usage Log.
Configuring Comet Memory#
Comet shares an off-heap memory pool with Spark. The size of the pool is
specified by spark.memory.offHeap.size. The pool is a shared budget rather than a shared
allocator: Comet’s native operators allocate from the Rust heap rather than from JVM off-heap
memory, but every reservation they make is charged against this same pool, so Comet and Spark’s own
off-heap consumers draw down one number.
A larger pool lets the operators that can spill, such as sorts, aggregations and the shuffle writer, keep more in memory before they spill to disk. An operator that cannot spill fails the task when it cannot reserve the memory it needs.
Comet’s memory pool only tracks memory that an operator explicitly reserves, which in practice means the batches an operator deliberately accumulates: the sort buffer, the build side of a hash join, hash aggregation state, and the shuffle writer’s buffered partitions. Memory that is not reserved is invisible to the pool no matter how much of it there is. That includes:
per-batch working memory in expression kernels and Arrow array builders,
decompression buffers and Parquet reader structures,
object store request buffers and the async runtime’s own machinery,
Arrow buffers allocated on the JVM side, which no budget covers at all,
allocator overhead: buffer padding, size-class rounding, fragmentation, and pages the allocator retains after a free rather than returning to the operating system.
Reserved memory is therefore a lower bound on what Comet really uses, and how far below it sits depends on the
workload. This is why Comet can stay within the pool’s limit and still push the executor past its container limit.
The part that is not counted has to fit in spark.executor.memoryOverhead, and each executor logs how large it is
while Comet runs; see Sizing the Overhead from the Memory Usage Log.
spark.comet.exec.memoryPool.fraction is deprecated and does not leave room for it. Spark hands out all of
spark.memory.offHeap.size to the tasks that ask for it, whatever the fraction. The fair_unified pool applies the
fraction to each task separately, where Spark’s own limit of an even share of the pool per running task is tighter
whenever more than one task is running, and the greedy_unified pool ignores it.
For more details about Spark off-heap memory mode, please refer to Spark documentation.
Comet implements multiple memory pool implementations. The type of pool can be specified with spark.comet.exec.memoryPool.
The valid pool types are:
fair_unified(default whenspark.memory.offHeap.enabled=trueis set)greedy_unified
Both pool types are shared by all the native plans in the same Spark task. A task can run more than one native plan at a time, for example the native operators on either side of a union or a coalesce. The shared pool ensures that their combined memory usage stays within the per-task limit.
The fair_unified pool prevents operators from using more than an even fraction of the available memory
(i.e. pool_size / num_consumers, where num_consumers counts the memory consumers registered by all of the task’s
native plans). This pool works best when you know beforehand
the query has multiple operators that will likely all need to spill. Sometimes it will cause spills even
when there is sufficient memory in order to leave enough memory for other operators.
The greedy_unified pool type implements a greedy first-come first-serve limit. This pool works well for queries that do not
need to spill or have a single spillable operator.
Configuring Executor Memory Overhead#
Enabling off-heap memory is not sufficient on its own. The cluster manager sizes the executor
container to include spark.memory.offHeap.size, so the memory that Comet’s operators reserve has
room, but nothing in the container’s size accounts for the memory that no pool tracks (see
Configuring Comet Memory). That memory has to fit in spark.executor.memoryOverhead, which the
JVM’s own non-heap memory, such as metaspace, code cache, thread stacks and GC structures, is
already drawing on.
Work out what the executor already gets before choosing a value. When
spark.executor.memoryOverhead is unset, Spark derives the overhead as
max(spark.executor.memoryOverheadFactor * spark.executor.memory, 384 MiB). The factor defaults to
0.1. On Kubernetes, when spark.executor.memoryOverheadFactor is unset, Spark uses
spark.kubernetes.memoryOverheadFactor instead, which also defaults to 0.1, except for PySpark
and SparkR applications submitted in cluster mode, where it defaults to 0.4. On Spark 4.0 and
later the floor is configurable through spark.executor.minMemoryOverhead. Setting
spark.executor.memoryOverhead replaces the derived value rather than adding to it, so a value
below what is derived today shrinks the container instead of growing it.
For a small executor, 2g is a reasonable starting point. A 4 GiB executor derives only 409 MiB, so
this is a real increase:
spark.executor.memoryOverhead=2g
A 32 GiB executor, on the other hand, already derives 3276 MiB, and the same setting would take away
1228 MiB. For executors that large, either pick an absolute value above what is derived today, or
raise spark.executor.memoryOverheadFactor instead so that the overhead keeps scaling with executor
size:
spark.executor.memoryOverheadFactor=0.2
Raise the value further if executors are killed by the cluster manager (on Kubernetes,
ExecutorLostFailure with exit code 137) rather than failing with a task-level out-of-memory error.
To measure how much Comet needs rather than guessing, see Sizing the Overhead from the Memory Usage Log.
Note that on Kubernetes and YARN the overhead is added to the container size, so raising it reduces
how many executors fit on a node. A standalone cluster ignores the overhead: its workers start
executors without a memory limit and count only spark.executor.memory against the memory they
offer.
Sizing the Overhead from the Memory Usage Log#
While Comet native plans are running, each executor logs its native memory usage at INFO level, one line every 10 seconds for the whole executor:
Comet native memory usage: allocated 5412.3 MiB, reserved 3890.0 MiB (16 native plans, 8 memory pools); JVM Arrow allocated 310.4 MiB, 96.2 MiB of it imported from native
allocatedis the memory that Comet’s native code has allocated and not yet freed, whether or not a pool tracks it.reservedis the part that Comet’s memory pools have reserved from Spark’s off-heap memory. It is charged againstspark.memory.offHeap.size, so the container already has room for it. A pool sometimes has to track memory that Spark could not grant, such as a spilled batch read back from disk while the off-heap memory is full.reservedleaves that memory out, since nothing charges it againstspark.memory.offHeap.size.JVM Arrow allocatedis the Arrow memory Comet holds on the JVM side, such as batches read from Comet’s in-memory cache, broadcast data, and batches exchanged with native code or Python workers. The part imported from native was allocated by Comet’s native code, soallocatedalready counts it, and the rest the JVM allocated itself. A pool reserves the JVM’s part only while a native operator holds on to one of its batches, such as a sort buffering its input. The split follows which allocator a buffer is charged to rather than where it was allocated, so treat the JVM’s part as an estimate.
When the application writes an event log and runs the Comet plugin, the event log also records a summary of these samples; see Reading the Memory Usage Log from the Event Log.
Comet’s untracked memory is what Comet holds outside the JVM heap that Spark’s off-heap pool does
not account for: allocated, plus the JVM Arrow figure less the part imported from native, minus
reserved. It is the part of Comet’s footprint that has to fit in spark.executor.memoryOverhead,
alongside the JVM’s own non-heap memory. To size the overhead from it:
Run a representative workload and find the line with the most untracked memory in each executor’s log, or the event with the most for each executor in the event log. Take all the figures from the same line or event: they are sampled together, and figures from different samples describe different moments. Setting
spark.comet.memory.logInterval=1sfor this run makes a short-lived peak less likely to fall between samples.Start from the overhead the executors had before Comet was enabled, which covers the JVM’s own non-heap memory, and add the most untracked memory seen on any executor.
Add a margin on top. The log can miss the true peak between samples, and none of the figures includes the allocator’s fragmentation and retained pages, or memory allocated by native C libraries such as zstd.
For example, a 16 GiB executor derives an overhead of 1638 MiB. If the line above has the most
untracked memory in its log, that is 5412.3 + (310.4 - 96.2) - 3890.0 = 1736.5 MiB. The overhead
then needs to be at least 1638 + 1737 = 3375 MiB before any margin, so
spark.executor.memoryOverhead=4g would be a reasonable setting.
The executor also logs a warning when its native memory looks larger than its container allows:
when the untracked memory, plus everything in use in Spark’s off-heap memory pool (which includes
Comet’s reservations), exceeds spark.memory.offHeap.size plus the memory overhead. This counts the
part of the off-heap pool that nothing has acquired at that moment, which untracked memory can
occupy until Spark hands it out, so a quiet log is not a sign that the overhead is large enough:
size it from the most untracked memory as described above. The overhead also has to hold the JVM’s
own non-heap memory, so by the time the warning appears the executor has likely outgrown its
container. It warns the first time this happens, and again each time it happens after dropping back
below. It derives the overhead the way YARN and Kubernetes size the executor’s container, as
described in Configuring Executor Memory Overhead. For a PySpark application it also counts
spark.executor.pyspark.memory as part of the container when the cluster manager adds it, which
YARN does, and Kubernetes does in cluster mode. There is no warning in local mode, on a standalone
cluster, whose workers do not limit an executor’s memory, or with any other cluster manager.
Look more closely before raising the overhead if the untracked memory keeps growing through a run
rather than levelling off: memory that is not being released will exhaust any overhead eventually.
The executor logs one more line after its last native plan finishes, and an allocated or
JVM Arrow allocated figure there that grows from one query to the next points the same way.
spark.comet.memory.logInterval is read when an executor starts its first Comet native plan, so set
it when the application is submitted. Set it to 0 to turn the log off.
Reading the Memory Usage Log from the Event Log#
When spark.eventLog.enabled is true and the application runs the Comet plugin
(spark.plugins=org.apache.spark.CometPlugin), each executor also sends its samples to the driver,
which writes a summary of them to the application’s event log. The event log keeps every executor’s
samples in one place, next to the jobs and stages they ran alongside, while executor logs are spread
over the cluster and often go away with their containers.
The driver writes, and flushes, each event of its event log as it goes, so it does not write every
sample. For each executor, it takes the samples a minute at a time and writes two of them: the one
with the most untracked memory, and the last, which in the minute an executor goes idle is the one
after its last native plan finishes. It writes them at the executor’s first heartbeat a minute or
more after the first of them arrived, which comes every spark.executor.heartbeatInterval, 10
seconds by default, busy or not, and when the executor is removed, such as when the cluster manager
kills it, or the application stops. So each executor adds about two events a minute while it runs
native plans, however short the interval, and each minute’s two reach the event log within about a
minute of that minute’s first sample arriving. The event with the most untracked memory for an
executor is then its sample with the most, apart from samples it takes while the application stops,
and samples from the last minute before a driver exits without stopping the application.
Each sample the driver writes is one event, on one line of the event log. It carries the figures of the line in bytes, the executor that took the sample, and when it did, in milliseconds since the epoch by the executor’s clock. The line above looks like this as an event, spread over several lines here:
{
"Event": "org.apache.comet.CometExecutorMemoryUsage",
"executorId": "3",
"time": 1727550000000,
"nativeAllocated": 5675212390,
"poolsReserved": 4078960640,
"pools": 8,
"plans": 16,
"jvmArrowAllocated": 325478809,
"jvmArrowImported": 100873420
}
nativeAllocated is allocated, poolsReserved is reserved, and jvmArrowAllocated and
jvmArrowImported are the JVM Arrow figures. With jq, for example, this lists the sample with the
most untracked memory for each executor, in MiB:
jq -n -c '[inputs | select(.Event == "org.apache.comet.CometExecutorMemoryUsage")
| {executorId, time, untrackedMiB: (([.nativeAllocated - .poolsReserved
+ ([.jvmArrowAllocated - .jvmArrowImported, 0] | max), 0] | max) / 1048576 | ceil)}]
| group_by(.executorId) | map(max_by(.untrackedMiB))[]' <event log files>
From Spark 4.0, event logs are compressed with zstd and split into several files by default, so
decompress the files first, for example with zstd -dc. The Spark history server does not display
these events, and one without Comet on its classpath logs once that it dropped them. From Spark 4.1,
spark.eventLog.excludedPatterns=org.apache.comet.CometExecutorMemoryUsage leaves them out of the
event log.
Batch Size#
Comet processes data in columnar batches. The batch size is controlled by spark.comet.batchSize (default
8192 rows). Larger batches generally improve throughput by amortizing per-batch overhead, but they also
increase peak memory usage — a batch holds all projected columns in Arrow format at once. Reduce this value
if you see frequent spilling or out-of-memory errors on wide tables; increase it (for example to 16384) on
narrow tables when memory is plentiful.
spark.comet.shuffle.jvm.batchSize controls the batch size used when the JVM columnar shuffle writer
flushes sorted spill files. A value larger than spark.comet.batchSize is capped at spark.comet.batchSize.
Limiting Spill Disk Usage#
Native operators that spill to disk (aggregate, sort, shuffle) are bounded by
spark.comet.maxTempDirectorySize (default 100 GB). The operators of one Comet native plan share
the limit. A Spark task can run more than one native plan at a time, for example the native
operators on either side of a union or a coalesce, so an executor running N concurrent tasks may
use more than N times this value on shared local disks. If the limit is reached, further spills
fail and the query errors out. Raise this on workloads with large sort/aggregate/shuffle spills, or
lower it to protect executors on shared disks, remembering that the total across an executor is a
multiple of this value.