Skip to content

feat: write the executor memory usage log to the Spark event log - #6401

Merged
andygrove merged 5 commits into
apache:mainfrom
comphead:memory-usage-event-log
Oct 1, 2026
Merged

andygrove merged 5 commits into
apache:mainfrom
comphead:memory-usage-event-log

Conversation

@comphead

@comphead comphead commented Sep 29, 2026 •

Copy link
Copy Markdown
Contributor

Which issue does this PR close?

There is no issue for this change.

Rationale for this change

Each executor logs its native memory usage every spark.comet.memory.logInterval, 10 seconds by default, and the memory tuning guide sizes spark.executor.memoryOverhead from the most untracked memory any executor logged during a representative run. Executor logs are spread over the cluster, are collected differently on YARN and Kubernetes, and often go away with their containers, so finding that peak means gathering every executor's log first. The event log already keeps the rest of the application's history in one place, next to the jobs and stages the samples ran alongside, and it outlives the application.

What changes are included in this PR?

  • CometExecutorMemoryUsage is a new SparkListenerEvent that carries one sample: the figures of the log line in bytes (nativeAllocated, poolsReserved, pools, plans, jvmArrowAllocated, jvmArrowImported), the executorId, and the executor's time. EventLoggingListener writes an event it has no format of its own for with Jackson, as one JSON line with "Event" set to the class name. The history server does not display it. One without Comet on its classpath logs once that it dropped it.
  • CometExecIterator: when spark.eventLog.enabled is true and the executor runs the Comet plugin, each sample also goes to the driver. The executor keeps logging every sample at INFO, so the event log only adds to the executor's log. Executors receive the driver's spark.* settings, so the executor plugin reads whether the driver writes an event log, through the same config entry the driver reads, which ignores surrounding whitespace. A send that fails on the executor, such as when the driver endpoint cannot be found, stops the sending after one warning.
  • CometDriverPlugin keeps, for each executor, the samples that the event log has yet to record, and writes two of each minute of them: the one with the most untracked memory, and the last. A listener on a listener bus queue of its own writes them at the executor's first heartbeat a minute or more after the driver received the first of them. Executors send heartbeats every spark.executor.heartbeatInterval, 10 seconds by default, whether or not they are busy, so an idle executor's last minute reaches the event log too, and the event log does not record the heartbeats themselves. The minute is by the driver's clock, since an executor's can differ. The listener also writes what the driver holds of an executor's samples when the executor is removed, and of every executor's when the application ends. From Spark 4.0 the plugin's shutdown(), which Spark calls before it stops the listener bus, writes what is left as well. The driver holds the samples rather than the executor, because an executor cannot send what it holds when the application stops: stopExecutors() sends StopExecutor one-way, and the driver stops its listener bus shortly after. The listener has a queue of its own because on Spark 3.4 and 3.5 the listener bus stops before the plugins shut down, and the shared queue also runs spark.extraListeners and query execution listeners, which could hold the application's end back until the bus has stopped. receive returns null, because the message is one-way.
  • CometExecutorPlugin keeps its PluginContext for this and clears it at shutdown, unless a later plugin in the same JVM has already replaced it.
  • Docs: the memory tuning guide gains a section on reading the samples from the event log, with a jq recipe that lists each executor's peak untracked memory, and notes that spark.eventLog.excludedPatterns leaves the events out from Spark 4.1. The config description, the plugin overview and the memory management guide mention the new destination.

EventLoggingListener flushes the log for every event of this kind, so the driver writes about two events per busy executor per minute, however short the interval, rather than one per interval. Spark's own executor metrics reach the event log as peaks in the same way, per task and, behind spark.eventLog.logStageExecutorMetrics, per stage. The events the driver writes when the application ends follow SparkListenerApplicationEnd, which the history server allows for.

How are these changes tested?

  • CometPluginsEventLogSuite (new) starts an application with the plugin and an uncompressed event log for each test. One sends a sample through the executor plugin, posts the end of the application to the listener bus, and checks that the driver writes the sample to the event log once, before the application actually stops. That pins the listener that Spark 3.4 and 3.5 depend on, even on the Spark 4.1 profile the pull request tier runs. Another checks that the driver writes what an executor sent when that executor is removed, once, before the application ends. The third runs the driver plugin on a ManualClock, moves it on a minute, and checks that the next heartbeat from an idle executor writes its samples.
  • CometPluginsSuite gains a test that the executor sends nothing when it runs the plugin but the application writes no event log.
  • CometExecIteratorLifecycleSuite gains a test that the event's fields map from Native.getMemoryUsage and the JVM Arrow figures and round-trip through JsonProtocol, which the event log and the history server both use. It also gains a test of the summary: that a minute ends a minute after the driver received the first sample, by its own clock whatever the executor's says, which two samples it keeps, and that untracked memory counts the Arrow memory the JVM allocated but not the part imported from native.

The suites run in local mode, where the executor plugin reaches the driver plugin inside one JVM. Across JVMs, Spark's plugin RPC Java-serializes the sample, which a case class supports, but no test covers that path. The pull request tier runs the suites against Spark 4.1 only, and the nightly run covers the other profiles.

@github-actions github-actions Bot added the enhancement New feature or request label Sep 29, 2026
@andygrove
andygrove self-requested a review September 29, 2026 15:22

@andygrove andygrove left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

I ran this in local-cluster[1,1,1024], so the executor was a separate JVM, which is the case the description says nothing covers. 27 samples from executor 0 went through the plugin RPC and read back from the event log with JsonProtocol, and the executor logged none of them at INFO. With extensions only, and with the plugin but no event log, the executor logged 27 INFO lines each and nothing reached the event log. The jq recipe also gives the guide's 1737 MiB for its example event. So the mechanism works. My two concerns are about how often it writes to the event log and what the executor keeps.

private[apache] def sendToEventLog(usage: Array[Long], jvmArrow: JvmArrowMemory): Boolean =
(Option(SparkEnv.get), CometExecutorPlugin.pluginContext) match {
case (Some(env), Some(pluginContext))
if !eventLogSendFailed && env.conf.getBoolean(EVENT_LOG_ENABLED, false) =>

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

This turns the event log path on for every application that runs the plugin and writes an event log, not only a sizing run, and I'm worried about what that costs the driver. EventLoggingListener sees these in onOtherEvent, which writes each one with flushLogger = true, so the single eventLog queue thread does an hflush per executor per interval. At the 1s interval the tuning guide recommends for a sizing run, a few hundred executors means a few hundred flushes a second. If that queue falls behind, AsyncEventQueue drops events from it, task and stage events included. Event log compaction also keeps every one of these lines, because none of Spark's event filters claim them.

Could the executor aggregate the samples and send a summary less often? It could keep sampling at the log interval, remember the sample with the most untracked memory since the last summary along with the latest one, and send them when the last plan finishes, when the container warning first fires, and otherwise at most once a minute or so. The sizing recipe only reads each executor's peak sample, so it gets the same answer, and the event count would depend on how long executors are busy rather than on the interval. That is close to what Spark does for its own executor metrics, which reach the event log only as per-stage peaks behind spark.eventLog.logStageExecutorMetrics, off by default. Sending on the warning also gets an executor's peak to the driver before it outgrows its container.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

Agreed on the cost. I moved the summary to the driver rather than the executor, though. An executor that holds a summary loses it when the application stops: stopExecutors() sends StopExecutor one-way and the driver stops its listener bus milliseconds later, so a run shorter than the window would keep only its first sample. Executors still send each sample. The driver keeps, per executor, the peak and the last sample of each minute, and writes what it holds when an executor is removed and when the application ends. That is about two events per busy executor per minute whatever the interval. The listener that writes at application end runs on a listener bus queue of its own, because on 3.4 and 3.5 the bus stops before the plugins shut down, and a slow listener on the shared queue could otherwise hold the application's end back until the bus has stopped.

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

PR description still needs updating?

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

Updated the description.

memoryUsageMessage(usage, jvmArrow, plansAtLastMemoryUsageLog).foreach(logInfo(_))
memoryUsageMessage(usage, jvmArrow, plansAtLastMemoryUsageLog).foreach { message =>
// A sample the event log records stays in the executor's log at DEBUG only.
if (sendToEventLog(usage, jvmArrow)) logDebug(message) else logInfo(message)

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

Could the executor keep this at INFO even when the sample goes to the event log? The send is fire-and-forget. For a remote driver, NettyRpcEnv.send only serializes the message and queues it on the outbox, so a connection failure lands in OneWayOutboxMessage.onFailure and never reaches the catch in sendToEventLog. On the driver, AsyncEventQueue.post drops the event when the eventLog queue is full. On Spark 4.1 and later, spark.eventLog.excludedPatterns can also filter these out. I tried that in local-cluster with the pattern set to org.apache.comet.CometExecutorMemoryUsage. 27 samples reached the driver's listener bus, none were written to the event log, and the executor logged none at INFO, so they were gone from both places. That setting is the obvious thing to reach for if these events make someone's event log too big.

One line per interval is what 1.1.0 already logs, so keeping it costs nothing new and makes the event log purely additive. It would also remove the DEBUG/INFO choice here, which is the one part of the change the new tests don't reach, since they call sendToEventLog directly.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

Done. Every sample is logged at INFO again, so the event log only adds to the executor's log. A send that fails outright now only stops the event log path, after one warning. The executor also reads spark.eventLog.enabled through Spark's typed entry, as the driver does, so a value with surrounding whitespace can no longer stop the log.

@comphead
comphead force-pushed the memory-usage-event-log branch from aee91d9 to 7b8ecae Compare September 30, 2026 00:32
@comphead comphead added the run-all-spark-profiles Run the Comet test suites against every Spark profile on this pull request, ahead of the merge queue label Sep 30, 2026
@comphead
comphead requested a review from andygrove September 30, 2026 01:00

@andygrove andygrove left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

Thanks for turning this around so quickly. Both of my comments are addressed. Every sample is logged at INFO again, and the event log gets about two events per busy executor per minute however short the interval. Your reason for holding the summary on the driver holds up. DriverEndpoint handles StopExecutors by sending StopExecutor one-way and replying straight away, so a summary held on the executor would be lost at stop. I also checked the stop order you rely on in the 3.4.3, 3.5.9, 4.0.4, 4.1.3 and 4.2.0 sources. On 3.x listenerBus.stop() comes before _plugins.foreach(_.shutdown()), and from 4.0 it comes after, so between them the listener and shutdown() cover every version.

The run-all-spark-profiles label run failed at startup without running any jobs. That's the problem #6409 fixes, so 3.4, 3.5, 4.0 and 4.2 haven't run on this yet. The branch also conflicts with main now, with one import line each in Plugins.scala and CometExecIteratorLifecycleSuite.scala. Merging main should take care of both, because the push runs the gated profiles now that the label is on the pull request.

I have one more concern, inline.

memoryUsageSummaries.remove(event.executorId).toList.flatMap(_.flush())
})

override def onApplicationEnd(event: SparkListenerApplicationEnd): Unit =

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

An executor that goes idle keeps its last minute of samples on the driver until it runs native plans again, is removed, or the application stops. On a long-lived application, like a Thrift or Connect server or a notebook session with static executors, that can be hours. Until then the event log doesn't have them. Someone who reads the in-progress log to size the overhead after a representative run misses each executor's last minute, which can be where the peak is. And if the driver dies before then, the samples are never written. My earlier ask had the summary go out when the last plan finishes for this reason.

Could this listener also end a summary that has been open for a minute, from onExecutorMetricsUpdate? Every executor heartbeat posts SparkListenerExecutorMetricsUpdate to the bus, every 10 seconds by default whether or not the executor is busy, and the event log doesn't write it. That would bound the wait to about a minute with no timer and no change to the event rate. The summary would need to note when the driver received its first sample, since start is by the executor's clock. It would also make the guide's "next to the jobs and stages they ran alongside" hold for an executor's last minute of work.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

Done. The listener now ends a summary from onExecutorMetricsUpdate once the driver received its first sample a minute or more before. That is the only place a minute ends now, and it goes by the driver's clock, so the executor-clock start is gone and receive only adds samples. A test runs the driver plugin on a ManualClock, moves it on a minute, and checks that the next heartbeat from an idle executor writes its samples. The guide now says the event log has every sample within about a minute of its arrival.

When spark.eventLog.enabled is true and the application runs the Comet
plugin, each executor sends its periodic native memory usage sample to
the driver plugin, which posts it to the listener bus as a
CometExecutorMemoryUsage event. The event log then keeps every
executor's samples in one place, and the executor logs the line at
DEBUG only.

Without the plugin, or after a failed send, the line stays in the
executor's log at INFO as before. The warning about exceeding the
executor's container always goes to the executor's log.
Addresses the review of the event log path.

The executor logs every memory usage sample at INFO again, so the
event log only adds to the executor's log. It still sends each sample
to the driver.

The driver keeps, for each executor, the sample with the most
untracked memory and the last sample of each minute, and writes only
those. EventLoggingListener flushes every event of this kind, so the
event log now gets about two events per busy executor per minute
however short the interval, rather than one per interval.

The driver holds the samples rather than the executor, because
stopExecutors() sends StopExecutor one-way and the driver stops its
listener bus shortly after, so an executor cannot send what it holds
when the application stops. A listener on a listener bus queue of its
own writes what the driver holds when an executor is removed and when
the application ends. It has its own queue because on Spark 3.4 and
3.5 the bus stops before the plugins shut down, and a slow listener on
the shared queue could otherwise hold the application's end back until
the bus has stopped. From Spark 4.0 the plugin's shutdown writes what
is left as well.

The executor plugin reads spark.eventLog.enabled through Spark's typed
entry, as the driver does, so a value with surrounding whitespace no
longer stops the memory usage log. The tuning guide's jq recipe clamps
untracked memory at zero, as the code does.
Addresses the review of the idle executor case.

An executor that went idle kept its last minute of memory usage samples
on the driver until it ran native plans again, was removed, or the
application stopped, which on a long-lived application can be hours.

The driver plugin's listener now ends an executor's summary at the
executor's first heartbeat a minute after the driver received the first
sample of it. Executors send a heartbeat every
spark.executor.heartbeatInterval whether or not they are busy, and the
event log does not record them. That is now the only place a minute
ends, by the driver's monotonic clock, so receive only adds samples and
the executor's clock no longer matters.

The driver plugin takes a Spark Clock, as HeartbeatReceiver does, so
that a test can move a ManualClock on a minute.
@comphead
comphead force-pushed the memory-usage-event-log branch from 7b8ecae to b1a1680 Compare September 30, 2026 15:31
@comphead
comphead requested a review from andygrove September 30, 2026 15:50

@andygrove andygrove left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

Thanks, this covers the idle case. I ran it in local-cluster[1,1,1024] on 4.1 and 3.4 with an executor that ran native plans for about 85 seconds and then sat idle with the application still up. Its last minute reached the event log about 54 seconds into the idle period, and the event with the most untracked memory matched the executor's INFO lines. The three suites pass on both versions, and the new test fails if the heartbeat doesn't end the summary. CI for this push is still queued behind other runs. I have one doc comment inline.

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 the event log has every sample within about a minute

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

The event log doesn't get every sample, only each minute's peak and last, as the start of this paragraph says. Could this say that each minute's samples reach the event log within about a minute of the first of them arriving?

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

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

Done. It now says each minute's two reach the event log within about a minute of that minute's first sample arriving.

The tuning guide said the event log has every sample within about a
minute of its arrival, but it records only each minute's peak and last.
It now says those two reach the event log within about a minute of that
minute's first sample arriving.
@andygrove
andygrove enabled auto-merge September 30, 2026 17:47
The event log tests asserted that the driver plugin's listener records a
sample after the event that makes it record it. The listener bus hands an
event to one queue at a time, and the Comet queue comes before the event
log's, so the listener's samples can reach the event log ahead of that
event. Spark 4.2's exec job failed that way.

Each test now posts a marker once every queue has handled the event, and
asserts that the event log records the sample exactly once, ahead of the
marker. The context posts its own application end only after the marker.
@andygrove
andygrove added this pull request to the merge queue Sep 30, 2026
Merged via the queue into apache:main with commit ec66afe Oct 1, 2026
55 checks passed
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

enhancement New feature or request run-all-spark-profiles Run the Comet test suites against every Spark profile on this pull request, ahead of the merge queue

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants