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]