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.

Spark and Comet both use the JVM heap and share the off-heap memory pool, and the rest of Comet's native memory has to fit in the executor's memory overhead

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 when spark.memory.offHeap.enabled=true is 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
  • allocated is the memory that Comet’s native code has allocated and not yet freed, whether or not a pool tracks it.

  • reserved is the part that Comet’s memory pools have reserved from Spark’s off-heap memory. It is charged against spark.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. reserved leaves that memory out, since nothing charges it against spark.memory.offHeap.size.

  • JVM Arrow allocated is 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, so allocated already 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:

  1. 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=1s for this run makes a short-lived peak less likely to fall between samples.

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

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