Skip to content
Closed
Show file tree
Hide file tree
Changes from 2 commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -827,6 +827,7 @@ private[spark] class AppStatusListener(
event.blockUpdatedInfo.blockId match {
case block: RDDBlockId => updateRDDBlock(event, block)
case stream: StreamBlockId => updateStreamBlock(event, stream)
case broadcast: BroadcastBlockId => updateBroadcastBlock(event, broadcast)
case _ =>
}
}
Expand Down Expand Up @@ -995,6 +996,34 @@ private[spark] class AppStatusListener(
}
}

private def updateBroadcastBlock(event: SparkListenerBlockUpdated,
Copy link
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

multi line args start on the next line.

broadcast: BroadcastBlockId): Unit = {
val now = System.nanoTime()
val executorId = event.blockUpdatedInfo.blockManagerId.executorId
val storageLevel = event.blockUpdatedInfo.storageLevel

// Whether values are being added to or removed from the existing accounting.
val diskDelta = event.blockUpdatedInfo.diskSize * (if (storageLevel.useDisk) 1 else -1)
val memoryDelta = event.blockUpdatedInfo.memSize * (if (storageLevel.useMemory) 1 else -1)

// Function to apply a delta to a value, but ensure that it doesn't go negative.
def newValue(old: Long, delta: Long): Long = math.max(0, old + delta)
Copy link
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Already exists (addDeltaToValue).


val maybeExec = liveExecutors.get(executorId)
maybeExec.foreach { exec =>
Copy link
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Move the foreach to the caller. That avoids repeating the foreach, and you could wrap more logic that doesn't need to run when the executor is not found.

if (exec.hasMemoryInfo) {
Copy link
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

This block exists in a very similar form in two other places. Feels like time to have a helper method.

if (storageLevel.useOffHeap) {
exec.usedOffHeap = newValue(exec.usedOffHeap, memoryDelta)
} else {
exec.usedOnHeap = newValue(exec.usedOnHeap, memoryDelta)
}
}
exec.memoryUsed = newValue(exec.memoryUsed, memoryDelta)
exec.diskUsed = newValue(exec.diskUsed, diskDelta)
maybeUpdate(exec, now)
}
}

private def getOrCreateStage(info: StageInfo): LiveStage = {
val stage = liveStages.computeIfAbsent((info.stageId, info.attemptNumber),
new Function[(Int, Int), LiveStage]() {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -938,6 +938,24 @@ class AppStatusListenerSuite extends SparkFunSuite with BeforeAndAfter {
intercept[NoSuchElementException] {
check[StreamBlockData](stream1.name) { _ => () }
}

// Update a BroadcastBlock.
val broadcast1 = BroadcastBlockId(1L)
listener.onBlockUpdated(SparkListenerBlockUpdated(
BlockUpdatedInfo(bm1, broadcast1, level, 1L, 1L)))

check[ExecutorSummaryWrapper](bm1.executorId) { exec =>
assert(exec.info.memoryUsed === 1L)
assert(exec.info.diskUsed === 1L)
}

// Drop a BroadcastBlock.
listener.onBlockUpdated(SparkListenerBlockUpdated(
BlockUpdatedInfo(bm1, broadcast1, StorageLevel.NONE, 1L, 1L)))
check[ExecutorSummaryWrapper](bm1.executorId) { exec =>
assert(exec.info.memoryUsed === 0)
assert(exec.info.diskUsed === 0)
}
}

test("eviction of old data") {
Expand Down