Skip to content

Commit ac8d572

Browse files
committed
stage end
1 parent 3094c92 commit ac8d572

1 file changed

Lines changed: 3 additions & 1 deletion

File tree

client/src/main/scala/org/apache/celeborn/client/commit/ReducePartitionCommitHandler.scala

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -151,7 +151,9 @@ class ReducePartitionCommitHandler(
151151
override def markShuffleDataLost(shuffleId: Int): Unit = {
152152
logWarning(s"Marking shuffle $shuffleId data as lost due to unknown/crashed worker.")
153153
dataLostShuffleSet.add(shuffleId)
154-
setStageEnd(shuffleId) // unblocks all pending GetReducerFileGroup waiters immediately
154+
if (!isStageEnd(shuffleId)) {
155+
setStageEnd(shuffleId)
156+
}
155157
}
156158

157159
override def isPartitionInProcess(shuffleId: Int, partitionId: Int): Boolean = {

0 commit comments

Comments
 (0)