andygrove commented on code in PR #6401:
URL: https://github.com/apache/datafusion-comet/pull/6401#discussion_r4135714440
##########
spark/src/main/scala/org/apache/comet/CometExecIterator.scala:
##########
@@ -527,6 +551,55 @@ object CometExecIterator extends Logging {
}
}
+ /**
+ * Sends a memory usage sample to the driver, which posts it to the listener
bus for the event
+ * log to record, and returns whether it did. It does when the application
writes an event log,
+ * which `spark.eventLog.enabled` shows on the executor as well as on the
driver, and the
+ * executor runs the Comet plugin, whose channel to the driver carries the
sample. Both are
+ * checked for every sample rather than once, because in local mode each
SparkContext starts its
+ * own executor plugin. A failure to send is not retried: those that surface
here, such as a
+ * driver endpoint that could not be found, persist, so the log falls back
to the executor's log
+ * for good.
+ */
+ 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) =>
Review Comment:
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.
##########
spark/src/main/scala/org/apache/comet/CometExecIterator.scala:
##########
@@ -508,7 +529,10 @@ object CometExecIterator extends Logging {
try {
val usage = nativeLib.getMemoryUsage()
val jvmArrow = JvmArrowMemory.current()
- 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)
Review Comment:
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.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]