Skip to content
Open
Show file tree
Hide file tree
Changes from all 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 @@ -28,13 +28,9 @@ private[spark] trait SchedulingAlgorithm {

private[spark] class FIFOSchedulingAlgorithm extends SchedulingAlgorithm {
override def comparator(s1: Schedulable, s2: Schedulable): Boolean = {
val priority1 = s1.priority
val priority2 = s2.priority
var res = math.signum(priority1 - priority2)
var res = Integer.compare(s1.priority, s2.priority)

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

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

Non-blocking (P2): Please add an in-repo regression test that fails against the old subtraction for both changed comparisons: priorities spanning Int.MinValue and Int.MaxValue, and equal-priority schedulables whose stage IDs span the same bounds. PoolSuite currently fixes priority at 0 and uses only small stage IDs, so it passes with the old comparator and would not catch a reintroduction of this scheduler-stalling failure.

@wangyum wangyum Sep 23, 2026 •

Copy link
Copy Markdown
Member Author

Choose a reason for hiding this comment

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

Added a test.

if (res == 0) {
val stageId1 = s1.stageId
val stageId2 = s2.stageId
res = math.signum(stageId1 - stageId2)
res = Integer.compare(s1.stageId, s2.stageId)
}
res < 0
}
Expand Down
53 changes: 50 additions & 3 deletions core/src/test/scala/org/apache/spark/scheduler/PoolSuite.scala
Original file line number Diff line number Diff line change
Expand Up @@ -39,12 +39,16 @@ class PoolSuite extends SparkFunSuite with LocalSparkContext {
val APP_NAME = "PoolSuite"
val TEST_POOL = "testPool"

def createTaskSetManager(stageId: Int, numTasks: Int, taskScheduler: TaskSchedulerImpl)
: TaskSetManager = {
def createTaskSetManager(
stageId: Int,
numTasks: Int,
taskScheduler: TaskSchedulerImpl,
priority: Int = 0,
stageAttemptId: Int = 0): TaskSetManager = {
val tasks = Array.tabulate[Task[_]](numTasks) { i =>
new FakeTask(stageId, i, Nil)
}
new TaskSetManager(taskScheduler, new TaskSet(tasks, stageId, 0, 0, null,
new TaskSetManager(taskScheduler, new TaskSet(tasks, stageId, stageAttemptId, priority, null,
ResourceProfile.DEFAULT_RESOURCE_PROFILE_ID, None), 0)
}

Expand Down Expand Up @@ -407,6 +411,49 @@ class PoolSuite extends SparkFunSuite with LocalSparkContext {
s"and '$expectedBatchPrefix'.\nCaptured logs:\n${logs.mkString("\n")}")
}

test("SPARK-59674: FIFO orders priorities and stage ids at Int bounds") {
sc = new SparkContext(LOCAL, APP_NAME)
val taskScheduler = new TaskSchedulerImpl(sc)
// This sequence makes getSortedTaskSetQueue throw IllegalArgumentException:
// Comparison method violates its general contract!
// when priority or stageId is ordered with signum(a - b).
val boundaryValues = Array(
0, Int.MaxValue, 0, 0, Int.MaxValue, Int.MinValue, Int.MinValue,
Int.MaxValue, Int.MaxValue, 0,
0, Int.MaxValue, Int.MaxValue, 0, 0, 0, Int.MaxValue, Int.MaxValue,
Int.MinValue, Int.MaxValue, 0, 0, 0, Int.MinValue, Int.MinValue, Int.MinValue,
Int.MinValue, Int.MaxValue, 0, 0, 0, Int.MinValue, Int.MinValue, Int.MinValue,
Int.MaxValue, 0, Int.MinValue, Int.MinValue, Int.MaxValue, Int.MaxValue,
Int.MinValue, Int.MinValue,
Int.MaxValue, 0, Int.MinValue, Int.MinValue, Int.MaxValue, Int.MaxValue,
Int.MinValue, Int.MinValue)
val expected = boundaryValues.sorted.toSeq

// Equal stage ids so the priority comparison is the only one that can overflow.
val priorityPool = new Pool("", FIFO, 0, 0)
boundaryValues.zipWithIndex.foreach { case (priority, attempt) =>
priorityPool.addSchedulable(createTaskSetManager(
stageId = 0,
numTasks = 1,
taskScheduler,
priority = priority,
stageAttemptId = attempt))
}
assert(priorityPool.getSortedTaskSetQueue.map(_.priority) === expected)

// Equal priorities so ordering falls through to stageId.
val stagePool = new Pool("", FIFO, 0, 0)
boundaryValues.zipWithIndex.foreach { case (stageId, attempt) =>
stagePool.addSchedulable(createTaskSetManager(
stageId,
numTasks = 1,
taskScheduler,
priority = 1,
stageAttemptId = attempt))
}
assert(stagePool.getSortedTaskSetQueue.map(_.stageId) === expected)
}

private def verifyPool(rootPool: Pool, poolName: String, expectedInitMinShare: Int,
expectedInitWeight: Int, expectedSchedulingMode: SchedulingMode): Unit = {
val selectedPool = rootPool.getSchedulableByName(poolName)
Expand Down