From 2f4490840cb796f89f627f1b9105f1658f32e456 Mon Sep 17 00:00:00 2001 From: Tom Magrino Date: Fri, 10 Jun 2016 14:11:30 -0700 Subject: [PATCH 1/5] Add links to executor logs from stage details page in UI. --- .../apache/spark/ui/jobs/ExecutorTable.scala | 9 ++++ .../org/apache/spark/ui/jobs/StagePage.scala | 42 +++++++++++++++---- .../org/apache/spark/ui/jobs/StagesTab.scala | 1 + 3 files changed, 45 insertions(+), 7 deletions(-) diff --git a/core/src/main/scala/org/apache/spark/ui/jobs/ExecutorTable.scala b/core/src/main/scala/org/apache/spark/ui/jobs/ExecutorTable.scala index 293f1438b8d57..5a97b174e443f 100644 --- a/core/src/main/scala/org/apache/spark/ui/jobs/ExecutorTable.scala +++ b/core/src/main/scala/org/apache/spark/ui/jobs/ExecutorTable.scala @@ -85,6 +85,7 @@ private[ui] class ExecutorTable(stageId: Int, stageAttemptId: Int, parent: Stage Shuffle Spill (Memory) Shuffle Spill (Disk) }} + Logs {createExecutorTable()} @@ -149,6 +150,14 @@ private[ui] class ExecutorTable(stageId: Int, stageAttemptId: Int, parent: Stage {Utils.bytesToString(v.diskBytesSpilled)} }} + + {val logs = parent.executorsListener.executorToLogUrls(k) + if (logs.isEmpty) { + "No Logs Found" + } else logs.map { + case (logName, logUrl) =>
{logName}
+ }} + } case None => diff --git a/core/src/main/scala/org/apache/spark/ui/jobs/StagePage.scala b/core/src/main/scala/org/apache/spark/ui/jobs/StagePage.scala index d986a55959b82..24207f75071ce 100644 --- a/core/src/main/scala/org/apache/spark/ui/jobs/StagePage.scala +++ b/core/src/main/scala/org/apache/spark/ui/jobs/StagePage.scala @@ -30,6 +30,7 @@ import org.apache.spark.SparkConf import org.apache.spark.executor.TaskMetrics import org.apache.spark.scheduler.{AccumulableInfo, TaskInfo, TaskLocality} import org.apache.spark.ui._ +import org.apache.spark.ui.exec.ExecutorsListener import org.apache.spark.ui.jobs.UIData._ import org.apache.spark.util.{Distribution, Utils} @@ -39,6 +40,7 @@ private[ui] class StagePage(parent: StagesTab) extends WebUIPage("stage") { private val progressListener = parent.progressListener private val operationGraphListener = parent.operationGraphListener + private val executorsListener = parent.executorsListener private val TIMELINE_LEGEND = {
@@ -296,7 +298,8 @@ private[ui] class StagePage(parent: StagesTab) extends WebUIPage("stage") { currentTime, pageSize = taskPageSize, sortColumn = taskSortColumn, - desc = taskSortDesc + desc = taskSortDesc, + executorsListener = executorsListener ) (_taskTable, _taskTable.table(page)) } catch { @@ -835,7 +838,8 @@ private[ui] class TaskTableRowData( val shuffleRead: Option[TaskTableRowShuffleReadData], val shuffleWrite: Option[TaskTableRowShuffleWriteData], val bytesSpilled: Option[TaskTableRowBytesSpilledData], - val error: String) + val error: String, + val logs: Map[String, String]) private[ui] class TaskDataSource( tasks: Seq[TaskUIData], @@ -848,7 +852,8 @@ private[ui] class TaskDataSource( currentTime: Long, pageSize: Int, sortColumn: String, - desc: Boolean) extends PagedDataSource[TaskTableRowData](pageSize) { + desc: Boolean, + executorsListener: ExecutorsListener) extends PagedDataSource[TaskTableRowData](pageSize) { import StagePage._ // Convert TaskUIData to TaskTableRowData which contains the final contents to show in the table @@ -992,6 +997,8 @@ private[ui] class TaskDataSource( None } + val logs = executorsListener.executorToLogUrls.getOrElse(info.executorId, Map.empty) + new TaskTableRowData( info.index, info.taskId, @@ -1015,7 +1022,8 @@ private[ui] class TaskDataSource( shuffleRead, shuffleWrite, bytesSpilled, - taskData.errorMessage.getOrElse("")) + taskData.errorMessage.getOrElse(""), + logs) } /** @@ -1193,6 +1201,16 @@ private[ui] class TaskDataSource( override def compare(x: TaskTableRowData, y: TaskTableRowData): Int = Ordering.String.compare(x.error, y.error) } + case "Logs" => new Ordering[TaskTableRowData] { + override def compare(x: TaskTableRowData, y: TaskTableRowData): Int = + if (x.logs.isEmpty == y.logs.isEmpty) { + return 0 + } else if (x.logs.isEmpty) { + return -1; + } else { + return 1; + } + } case unknownColumn => throw new IllegalArgumentException(s"Unknown column: $unknownColumn") } if (desc) { @@ -1217,7 +1235,8 @@ private[ui] class TaskPagedTable( currentTime: Long, pageSize: Int, sortColumn: String, - desc: Boolean) extends PagedTable[TaskTableRowData] { + desc: Boolean, + executorsListener: ExecutorsListener) extends PagedTable[TaskTableRowData] { // We only track peak memory used for unsafe operators private val displayPeakExecutionMemory = conf.getBoolean("spark.sql.unsafe.enabled", true) @@ -1244,7 +1263,8 @@ private[ui] class TaskPagedTable( currentTime, pageSize, sortColumn, - desc) + desc, + executorsListener) override def pageLink(page: Int): String = { val encodedSortColumn = URLEncoder.encode(sortColumn, "UTF-8") @@ -1297,7 +1317,8 @@ private[ui] class TaskPagedTable( } else { Nil }} ++ - Seq(("Errors", "")) + Seq(("Errors", "")) ++ + Seq(("Logs", "")) if (!taskHeadersAndCssClasses.map(_._1).contains(sortColumn)) { throw new IllegalArgumentException(s"Unknown column: $sortColumn") @@ -1391,6 +1412,13 @@ private[ui] class TaskPagedTable( {task.bytesSpilled.get.diskBytesSpilledReadable} }} {errorMessageCell(task.error)} + + {if (task.logs.isEmpty) { + "No Logs Found" + } else task.logs.map { + case (logName, logUrl) =>
{logName}
+ }} + } diff --git a/core/src/main/scala/org/apache/spark/ui/jobs/StagesTab.scala b/core/src/main/scala/org/apache/spark/ui/jobs/StagesTab.scala index bd5f16d25b477..573192ac17d45 100644 --- a/core/src/main/scala/org/apache/spark/ui/jobs/StagesTab.scala +++ b/core/src/main/scala/org/apache/spark/ui/jobs/StagesTab.scala @@ -29,6 +29,7 @@ private[ui] class StagesTab(parent: SparkUI) extends SparkUITab(parent, "stages" val killEnabled = parent.killEnabled val progressListener = parent.jobProgressListener val operationGraphListener = parent.operationGraphListener + val executorsListener = parent.executorsListener attachPage(new AllStagesPage(this)) attachPage(new StagePage(this)) From 482175294010a01cdd21c9ff60068b0fcc9180a5 Mon Sep 17 00:00:00 2001 From: Tom Magrino Date: Fri, 10 Jun 2016 16:59:08 -0700 Subject: [PATCH 2/5] Small fixs for missing indices and testing. --- .../scala/org/apache/spark/ui/jobs/ExecutorTable.scala | 2 +- .../test/scala/org/apache/spark/ui/StagePageSuite.scala | 8 +++++--- 2 files changed, 6 insertions(+), 4 deletions(-) diff --git a/core/src/main/scala/org/apache/spark/ui/jobs/ExecutorTable.scala b/core/src/main/scala/org/apache/spark/ui/jobs/ExecutorTable.scala index 5a97b174e443f..1c2c36b477038 100644 --- a/core/src/main/scala/org/apache/spark/ui/jobs/ExecutorTable.scala +++ b/core/src/main/scala/org/apache/spark/ui/jobs/ExecutorTable.scala @@ -151,7 +151,7 @@ private[ui] class ExecutorTable(stageId: Int, stageAttemptId: Int, parent: Stage }} - {val logs = parent.executorsListener.executorToLogUrls(k) + {val logs = parent.executorsListener.executorToLogUrls.getOrElse(k, Map.empty) if (logs.isEmpty) { "No Logs Found" } else logs.map { diff --git a/core/src/test/scala/org/apache/spark/ui/StagePageSuite.scala b/core/src/test/scala/org/apache/spark/ui/StagePageSuite.scala index 6d726d3d591b5..0fec0d959fdc7 100644 --- a/core/src/test/scala/org/apache/spark/ui/StagePageSuite.scala +++ b/core/src/test/scala/org/apache/spark/ui/StagePageSuite.scala @@ -20,12 +20,12 @@ package org.apache.spark.ui import javax.servlet.http.HttpServletRequest import scala.xml.Node - -import org.mockito.Mockito.{mock, when, RETURNS_SMART_NULLS} - +import org.mockito.Mockito.{RETURNS_SMART_NULLS, mock, when} import org.apache.spark._ import org.apache.spark.executor.TaskMetrics import org.apache.spark.scheduler._ +import org.apache.spark.storage.StorageStatusListener +import org.apache.spark.ui.exec.ExecutorsListener import org.apache.spark.ui.jobs.{JobProgressListener, StagePage, StagesTab} import org.apache.spark.ui.scope.RDDOperationGraphListener @@ -64,11 +64,13 @@ class StagePageSuite extends SparkFunSuite with LocalSparkContext { private def renderStagePage(conf: SparkConf): Seq[Node] = { val jobListener = new JobProgressListener(conf) val graphListener = new RDDOperationGraphListener(conf) + val executorsListener = new ExecutorsListener(new StorageStatusListener(conf), conf) val tab = mock(classOf[StagesTab], RETURNS_SMART_NULLS) val request = mock(classOf[HttpServletRequest]) when(tab.conf).thenReturn(conf) when(tab.progressListener).thenReturn(jobListener) when(tab.operationGraphListener).thenReturn(graphListener) + when(tab.executorsListener).thenReturn(executorsListener) when(tab.appName).thenReturn("testing") when(tab.headerTabs).thenReturn(Seq.empty) when(request.getParameter("id")).thenReturn("0") From 0942220ee1cd7c94c36db2f75670eba44e94d506 Mon Sep 17 00:00:00 2001 From: Tom Magrino Date: Mon, 13 Jun 2016 09:17:59 -0700 Subject: [PATCH 3/5] Code style fix. --- core/src/test/scala/org/apache/spark/ui/StagePageSuite.scala | 4 +++- 1 file changed, 3 insertions(+), 1 deletion(-) diff --git a/core/src/test/scala/org/apache/spark/ui/StagePageSuite.scala b/core/src/test/scala/org/apache/spark/ui/StagePageSuite.scala index 0fec0d959fdc7..d30b987d6ca31 100644 --- a/core/src/test/scala/org/apache/spark/ui/StagePageSuite.scala +++ b/core/src/test/scala/org/apache/spark/ui/StagePageSuite.scala @@ -20,7 +20,9 @@ package org.apache.spark.ui import javax.servlet.http.HttpServletRequest import scala.xml.Node -import org.mockito.Mockito.{RETURNS_SMART_NULLS, mock, when} + +import org.mockito.Mockito.{mock, when, RETURNS_SMART_NULLS} + import org.apache.spark._ import org.apache.spark.executor.TaskMetrics import org.apache.spark.scheduler._ From 70d54e801e4b245c06fa111c72e221dedbd08607 Mon Sep 17 00:00:00 2001 From: Tom Magrino Date: Tue, 14 Jun 2016 12:23:24 -0700 Subject: [PATCH 4/5] Changes to address UI consistency comments by @ajbozarth. --- .../scala/org/apache/spark/ui/jobs/ExecutorTable.scala | 4 +--- .../main/scala/org/apache/spark/ui/jobs/StagePage.scala | 8 +++----- 2 files changed, 4 insertions(+), 8 deletions(-) diff --git a/core/src/main/scala/org/apache/spark/ui/jobs/ExecutorTable.scala b/core/src/main/scala/org/apache/spark/ui/jobs/ExecutorTable.scala index 1c2c36b477038..6371e0e107614 100644 --- a/core/src/main/scala/org/apache/spark/ui/jobs/ExecutorTable.scala +++ b/core/src/main/scala/org/apache/spark/ui/jobs/ExecutorTable.scala @@ -152,9 +152,7 @@ private[ui] class ExecutorTable(stageId: Int, stageAttemptId: Int, parent: Stage }} {val logs = parent.executorsListener.executorToLogUrls.getOrElse(k, Map.empty) - if (logs.isEmpty) { - "No Logs Found" - } else logs.map { + logs.map { case (logName, logUrl) => }} diff --git a/core/src/main/scala/org/apache/spark/ui/jobs/StagePage.scala b/core/src/main/scala/org/apache/spark/ui/jobs/StagePage.scala index 24207f75071ce..eaa60ef524a29 100644 --- a/core/src/main/scala/org/apache/spark/ui/jobs/StagePage.scala +++ b/core/src/main/scala/org/apache/spark/ui/jobs/StagePage.scala @@ -1201,7 +1201,7 @@ private[ui] class TaskDataSource( override def compare(x: TaskTableRowData, y: TaskTableRowData): Int = Ordering.String.compare(x.error, y.error) } - case "Logs" => new Ordering[TaskTableRowData] { + case "Executor Logs" => new Ordering[TaskTableRowData] { override def compare(x: TaskTableRowData, y: TaskTableRowData): Int = if (x.logs.isEmpty == y.logs.isEmpty) { return 0 @@ -1318,7 +1318,7 @@ private[ui] class TaskPagedTable( Nil }} ++ Seq(("Errors", "")) ++ - Seq(("Logs", "")) + Seq(("Executor Logs", "")) if (!taskHeadersAndCssClasses.map(_._1).contains(sortColumn)) { throw new IllegalArgumentException(s"Unknown column: $sortColumn") @@ -1413,9 +1413,7 @@ private[ui] class TaskPagedTable( }} {errorMessageCell(task.error)} - {if (task.logs.isEmpty) { - "No Logs Found" - } else task.logs.map { + {task.logs.map { case (logName, logUrl) => }} From 615e990b9a4e90b54a5e7ef28664070970d3c88b Mon Sep 17 00:00:00 2001 From: Tom Magrino Date: Fri, 1 Jul 2016 11:31:22 -0700 Subject: [PATCH 5/5] More minimalistic design as requested in PR feedback. --- .../apache/spark/ui/jobs/ExecutorTable.scala | 19 +++++++----- .../org/apache/spark/ui/jobs/StagePage.scala | 29 +++++++------------ 2 files changed, 22 insertions(+), 26 deletions(-) diff --git a/core/src/main/scala/org/apache/spark/ui/jobs/ExecutorTable.scala b/core/src/main/scala/org/apache/spark/ui/jobs/ExecutorTable.scala index 6371e0e107614..133c3b1b9aca8 100644 --- a/core/src/main/scala/org/apache/spark/ui/jobs/ExecutorTable.scala +++ b/core/src/main/scala/org/apache/spark/ui/jobs/ExecutorTable.scala @@ -85,7 +85,6 @@ private[ui] class ExecutorTable(stageId: Int, stageAttemptId: Int, parent: Stage Shuffle Spill (Memory) Shuffle Spill (Disk) }} - Logs {createExecutorTable()} @@ -115,7 +114,17 @@ private[ui] class ExecutorTable(stageId: Int, stageAttemptId: Int, parent: Stage case Some(stageData: StageUIData) => stageData.executorSummary.toSeq.sortBy(_._1).map { case (k, v) => - {k} + +
{k}
+
+ { + val logs = parent.executorsListener.executorToLogUrls.getOrElse(k, Map.empty) + logs.map { + case (logName, logUrl) => + } + } +
+ {executorIdToAddress.getOrElse(k, "CANNOT FIND ADDRESS")} {UIUtils.formatDuration(v.taskTime)} {v.failedTasks + v.succeededTasks + v.killedTasks} @@ -150,12 +159,6 @@ private[ui] class ExecutorTable(stageId: Int, stageAttemptId: Int, parent: Stage {Utils.bytesToString(v.diskBytesSpilled)} }} - - {val logs = parent.executorsListener.executorToLogUrls.getOrElse(k, Map.empty) - logs.map { - case (logName, logUrl) => - }} - } case None => diff --git a/core/src/main/scala/org/apache/spark/ui/jobs/StagePage.scala b/core/src/main/scala/org/apache/spark/ui/jobs/StagePage.scala index eaa60ef524a29..15d377fc3912d 100644 --- a/core/src/main/scala/org/apache/spark/ui/jobs/StagePage.scala +++ b/core/src/main/scala/org/apache/spark/ui/jobs/StagePage.scala @@ -1201,16 +1201,6 @@ private[ui] class TaskDataSource( override def compare(x: TaskTableRowData, y: TaskTableRowData): Int = Ordering.String.compare(x.error, y.error) } - case "Executor Logs" => new Ordering[TaskTableRowData] { - override def compare(x: TaskTableRowData, y: TaskTableRowData): Int = - if (x.logs.isEmpty == y.logs.isEmpty) { - return 0 - } else if (x.logs.isEmpty) { - return -1; - } else { - return 1; - } - } case unknownColumn => throw new IllegalArgumentException(s"Unknown column: $unknownColumn") } if (desc) { @@ -1317,8 +1307,7 @@ private[ui] class TaskPagedTable( } else { Nil }} ++ - Seq(("Errors", "")) ++ - Seq(("Executor Logs", "")) + Seq(("Errors", "")) if (!taskHeadersAndCssClasses.map(_._1).contains(sortColumn)) { throw new IllegalArgumentException(s"Unknown column: $sortColumn") @@ -1362,7 +1351,16 @@ private[ui] class TaskPagedTable( {if (task.speculative) s"${task.attempt} (speculative)" else task.attempt.toString} {task.status} {task.taskLocality} - {task.executorIdAndHost} + +
{task.executorIdAndHost}
+
+ { + task.logs.map { + case (logName, logUrl) => + } + } +
+ {UIUtils.formatDate(new Date(task.launchTime))} {task.formatDuration} @@ -1412,11 +1410,6 @@ private[ui] class TaskPagedTable( {task.bytesSpilled.get.diskBytesSpilledReadable} }} {errorMessageCell(task.error)} - - {task.logs.map { - case (logName, logUrl) => - }} - }