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


##########
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:
   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.
   



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