Skip to content
Closed
Changes from 1 commit
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,30 @@ private[spark] class AppStatusListener(
}
}

def updateBroadcastBlock(event: SparkListenerBlockUpdated, broadcast: BroadcastBlockId): Unit = {
Copy link
Member

Choose a reason for hiding this comment

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

def -> private def?

Copy link
Contributor Author

Choose a reason for hiding this comment

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

Thank you for your review! I have submitted a new change.

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


liveExecutors.get(executorId).foreach { exec =>
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)
}
}

private def getOrCreateStage(info: StageInfo): LiveStage = {
val stage = liveStages.computeIfAbsent((info.stageId, info.attemptNumber),
new Function[(Int, Int), LiveStage]() {
Expand Down