andygrove commented on code in PR #6401:
URL: https://github.com/apache/datafusion-comet/pull/6401#discussion_r4144960283


##########
spark/src/main/scala/org/apache/spark/Plugins.scala:
##########
@@ -49,9 +53,35 @@ import org.apache.comet.iceberg.IcebergWriteReportListener
  */
 class CometDriverPlugin extends DriverPlugin with Logging {
 
+  // Set by init, before Spark delivers any message, and read on the RPC 
thread that delivers them.
+  @volatile private var sparkContext: SparkContext = _
+
+  // By executor, the memory usage samples that the event log has yet to 
record. The RPC thread
+  // that delivers samples shares it with the threads that record what is left 
of them.
+  private val memoryUsageSummaries = mutable.HashMap.empty[String, 
MemoryUsageSummary]
+
   override def init(sc: SparkContext, pluginContext: PluginContext): 
ju.Map[String, String] = {
     logInfo("CometDriverPlugin init")
 
+    sparkContext = sc
+    if (sc.conf.get(EVENT_LOG_ENABLED)) {
+      // A queue of its own, so that a slow listener on the shared queue 
cannot hold the
+      // application's end back until the listener bus has stopped, which 
drops what is posted
+      // after it.
+      sc.listenerBus.addToQueue(
+        new SparkListener {
+          // An executor that has gone away sends no later sample to end its 
summary.
+          override def onExecutorRemoved(event: SparkListenerExecutorRemoved): 
Unit =
+            recordMemoryUsage(memoryUsageSummaries.synchronized {
+              
memoryUsageSummaries.remove(event.executorId).toList.flatMap(_.flush())
+            })
+
+          override def onApplicationEnd(event: SparkListenerApplicationEnd): 
Unit =

Review Comment:
   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.



-- 
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]

Reply via email to