From 4f708f1f69115101102ef172e08fd6da7327c358 Mon Sep 17 00:00:00 2001 From: Yuming Wang Date: Sun, 20 Sep 2026 22:16:59 +0800 Subject: [PATCH 1/3] Fix Comparison method violates its general contract in FIFOSchedulingAlgorithm --- .../org/apache/spark/scheduler/SchedulingAlgorithm.scala | 8 ++------ 1 file changed, 2 insertions(+), 6 deletions(-) diff --git a/core/src/main/scala/org/apache/spark/scheduler/SchedulingAlgorithm.scala b/core/src/main/scala/org/apache/spark/scheduler/SchedulingAlgorithm.scala index 18ebbbe78a5b4..a7db12dd97c58 100644 --- a/core/src/main/scala/org/apache/spark/scheduler/SchedulingAlgorithm.scala +++ b/core/src/main/scala/org/apache/spark/scheduler/SchedulingAlgorithm.scala @@ -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) 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 } From 671ca3a74501bb2ffe4738e2e9360d932a6bbf85 Mon Sep 17 00:00:00 2001 From: Yuming Wang Date: Wed, 23 Sep 2026 14:05:00 +0800 Subject: [PATCH 2/3] Add test --- .../apache/spark/scheduler/PoolSuite.scala | 54 +++++++++++++++++-- 1 file changed, 51 insertions(+), 3 deletions(-) diff --git a/core/src/test/scala/org/apache/spark/scheduler/PoolSuite.scala b/core/src/test/scala/org/apache/spark/scheduler/PoolSuite.scala index 37de8338ad905..a06ca9e360938 100644 --- a/core/src/test/scala/org/apache/spark/scheduler/PoolSuite.scala +++ b/core/src/test/scala/org/apache/spark/scheduler/PoolSuite.scala @@ -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) } @@ -407,6 +411,50 @@ 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) + // TimSort only detects the overflow while merging runs. Three values mis-order + // but do not throw. 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) From e83dc6badff9c16e72f865311c362359e1ad187c Mon Sep 17 00:00:00 2001 From: Yuming Wang Date: Wed, 23 Sep 2026 14:06:08 +0800 Subject: [PATCH 3/3] Add test --- .../test/scala/org/apache/spark/scheduler/PoolSuite.scala | 5 ++--- 1 file changed, 2 insertions(+), 3 deletions(-) diff --git a/core/src/test/scala/org/apache/spark/scheduler/PoolSuite.scala b/core/src/test/scala/org/apache/spark/scheduler/PoolSuite.scala index a06ca9e360938..1c18ac662faff 100644 --- a/core/src/test/scala/org/apache/spark/scheduler/PoolSuite.scala +++ b/core/src/test/scala/org/apache/spark/scheduler/PoolSuite.scala @@ -414,9 +414,8 @@ class PoolSuite extends SparkFunSuite with LocalSparkContext { test("SPARK-59674: FIFO orders priorities and stage ids at Int bounds") { sc = new SparkContext(LOCAL, APP_NAME) val taskScheduler = new TaskSchedulerImpl(sc) - // TimSort only detects the overflow while merging runs. Three values mis-order - // but do not throw. This sequence makes getSortedTaskSetQueue throw - // IllegalArgumentException: Comparison method violates its general contract! + // 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,