Restructure backend shutdown to reduce unhelpful churn - #150
GeorgeJahad wants to merge 3 commits into
Conversation
ArmadaExecutorAllocator runs ticks on a ScheduledExecutorService. When ArmadaClusterManagerBackend.stop() began, an in-flight tick could still submit a fresh executor job to Armada � and stop() also called super.stop() before cancelling pending Armada jobs, so the driver's CoarseGrainedScheduler endpoint was torn down while Armada was still free to start pods for already-submitted executors. Late-arriving pods then failed RegisterExecutor with RpcEndpointNotFoundException. Two changes fix both legs of the race: - Add a volatile flag on the backend, set at the very top of stop(). The allocator consults it in tryAllocateExecutors and submitExecutorJobs and becomes a no-op once shutdown begins. - Reorder stop() to cancel pending (not-yet-registered) executor jobs *before* super.stop(), so Armada never starts pods that would race against a dead RpcEnv. Already-running executors keep the existing grace-period flow so they can exit cleanly and report Succeeded. Adds unit tests covering both behaviours. Co-Authored-By: Claude Opus 4.7 (1M context) <noreply@anthropic.com> Signed-off-by: George Jahad <github@blackbirdsystems.net>
Signed-off-by: George Jahad <github@blackbirdsystems.net>
Signed-off-by: George Jahad <github@blackbirdsystems.net>
Greptile SummaryThis PR restructures
Confidence Score: 4/5The shutdown restructuring is logically sound and well-tested for the new paths; the main gap is the dropped property-test coverage for gang-cardinality allocation logic. The new stopping flag, removedExecutors dedup, and onDisconnected branching are all correctly implemented and covered by targeted tests. The one concern is that safeRemoveExecutor permanently gates an executor in removedExecutors before calling removeExecutor, so a transient failure silently prevents any retry — acceptable in shutdown but worth tracking. The more notable gap is the wholesale removal of the computeToAllocate table test, which was the guard for an existing gang-scheduling bug fix; the underlying logic is untouched but now has no regression coverage. ArmadaExecutorAllocatorSuite.scala — the computeToAllocate property tests were removed and should be restored. Important Files Changed
Sequence Diagram%%{init: {'theme': 'neutral'}}%%
sequenceDiagram
participant Spark as Spark Core
participant Backend as ArmadaClusterManagerBackend
participant Allocator as ArmadaExecutorAllocator
participant Endpoint as ArmadaDriverEndpoint (RpcEnv alive)
participant Executor as Executor Pod
participant Watcher as ArmadaEventWatcher
Note over Backend: stop() called
Backend->>Backend: "stopping = true"
Backend->>Allocator: stop() — no new submissions
Backend->>Backend: pre-mark getExecutorIds() in executorsPendingToRemove
Backend->>Endpoint: stopExecutors() → StopExecutor RPC to all registered executors
Note over Endpoint: RpcEnv stays alive during grace period
alt Late RegisterExecutor arrives during grace period
Executor->>Endpoint: RegisterExecutor(executorId)
Endpoint->>Executor: reply(true) + send(StopExecutor)
Note over Executor: exits 0 → Succeeded in Lookout
end
Executor-->>Backend: onDisconnected (executorsPendingToRemove contains id)
Backend->>Backend: safeRemoveExecutor (no-op if already removed)
Note over Backend: grace period elapses
Backend->>Watcher: continues running — captures Succeeded/Failed/Cancelled events
Watcher-->>Backend: onExecutorSucceeded / onExecutorFailed
Backend->>Backend: markTerminal + executorsPendingToRemove.put(true) + safeRemoveExecutor
Backend->>Backend: cancelExecutorJobs() — cancel remaining non-terminal Armada jobs
Backend->>Spark: super.stop() — tears down RpcEnv / driver endpoint
Backend->>Watcher: eventWatcher.stop()
%%{init: {'theme': 'base', 'themeVariables': {"darkMode": true, "background": "#0d1117", "primaryColor": "#21262d", "primaryTextColor": "#e6edf3", "primaryBorderColor": "#8b949e", "lineColor": "#8b949e", "textColor": "#e6edf3", "edgeLabelBackground": "#161b22", "actorBkg": "#21262d", "actorBorder": "#8b949e", "actorTextColor": "#e6edf3", "actorLineColor": "#8b949e", "signalColor": "#8b949e", "signalTextColor": "#e6edf3", "noteBkgColor": "#373320", "noteBorderColor": "#d4a72c", "noteTextColor": "#f0e6c0", "labelBoxBkgColor": "#21262d", "labelBoxBorderColor": "#8b949e", "labelTextColor": "#e6edf3", "loopTextColor": "#e6edf3", "activationBkgColor": "#30363d", "activationBorderColor": "#8b949e"}}}%%
sequenceDiagram
participant Spark as Spark Core
participant Backend as ArmadaClusterManagerBackend
participant Allocator as ArmadaExecutorAllocator
participant Endpoint as ArmadaDriverEndpoint (RpcEnv alive)
participant Executor as Executor Pod
participant Watcher as ArmadaEventWatcher
Note over Backend: stop() called
Backend->>Backend: "stopping = true"
Backend->>Allocator: stop() — no new submissions
Backend->>Backend: pre-mark getExecutorIds() in executorsPendingToRemove
Backend->>Endpoint: stopExecutors() → StopExecutor RPC to all registered executors
Note over Endpoint: RpcEnv stays alive during grace period
alt Late RegisterExecutor arrives during grace period
Executor->>Endpoint: RegisterExecutor(executorId)
Endpoint->>Executor: reply(true) + send(StopExecutor)
Note over Executor: exits 0 → Succeeded in Lookout
end
Executor-->>Backend: onDisconnected (executorsPendingToRemove contains id)
Backend->>Backend: safeRemoveExecutor (no-op if already removed)
Note over Backend: grace period elapses
Backend->>Watcher: continues running — captures Succeeded/Failed/Cancelled events
Watcher-->>Backend: onExecutorSucceeded / onExecutorFailed
Backend->>Backend: markTerminal + executorsPendingToRemove.put(true) + safeRemoveExecutor
Backend->>Backend: cancelExecutorJobs() — cancel remaining non-terminal Armada jobs
Backend->>Spark: super.stop() — tears down RpcEnv / driver endpoint
Backend->>Watcher: eventWatcher.stop()
|
| import org.mockito.Mockito._ | ||
| import org.scalatest.funsuite.AnyFunSuite | ||
| import org.scalatest.matchers.should.Matchers | ||
| import org.scalatest.prop.TableDrivenPropertyChecks | ||
|
|
||
| import org.apache.spark.SparkConf | ||
| import org.apache.spark.resource.ResourceProfile | ||
|
|
||
| class ArmadaExecutorAllocatorSuite | ||
| extends AnyFunSuite | ||
| with Matchers | ||
| with TableDrivenPropertyChecks { | ||
| class ArmadaExecutorAllocatorSuite extends AnyFunSuite with Matchers { | ||
|
|
||
| private def newAllocator(batchSize: Int): ArmadaExecutorAllocator = { | ||
| val backend = mock(classOf[ArmadaClusterManagerBackend]) | ||
| val armadaClient = mock(classOf[ArmadaClient]) | ||
| val conf = new SparkConf(false) | ||
| .set("spark.armada.allocation.batchSize", batchSize.toString) | ||
| private def newAllocator( | ||
| backend: ArmadaClusterManagerBackend, | ||
| armadaClient: ArmadaClient | ||
| ): ArmadaExecutorAllocator = | ||
| new ArmadaExecutorAllocator( | ||
| armadaClient, | ||
| "test-queue", | ||
| "test-jobset", | ||
| conf, | ||
| new SparkConf(false), | ||
| "test-app", | ||
| backend | ||
| ) | ||
|
|
||
| // Pre-load totalExpectedExecutors so the foreach in tryAllocateExecutors actually iterates — | ||
| // otherwise the assertion below holds vacuously whether or not the isStopping gate exists. | ||
| private def primeDemand(allocator: ArmadaExecutorAllocator): Unit = { | ||
| val rp = ResourceProfile.getOrCreateDefaultProfile(new SparkConf(false)) | ||
| allocator.setTotalExpectedExecutors(Map(rp -> 4)) | ||
| } | ||
|
|
||
| test("computeToAllocate caps the batch by gang cardinality") { | ||
| val allocator = newAllocator(batchSize = 4) | ||
|
|
||
| // (gap, gangCardinality, expectedToAllocate) | ||
| val cases = Table( | ||
| ("gap", "gangCardinality", "expected"), | ||
| // gang smaller than batch -> capped to the gang size (the bug this fix targets) | ||
| (4, 1, 1), | ||
| (4, 2, 2), | ||
| // gang larger than batch -> batch size wins | ||
| (4, 10, 4), | ||
| // gang equals batch -> batch size | ||
| (4, 4, 4), | ||
| // no gang constraint (<= 0) -> fall back to batch size | ||
| (4, 0, 4), | ||
| (4, -1, 4), | ||
| // gap smaller than the effective batch -> gap wins | ||
| (2, 4, 2), | ||
| (1, 10, 1), | ||
| (3, 5, 3), | ||
| // gap = 0 -> never submit, regardless of cardinality | ||
| (0, 5, 0) | ||
| ) | ||
| test("tryAllocateExecutors is a no-op when backend is stopping") { | ||
| val backend = mock(classOf[ArmadaClusterManagerBackend]) | ||
| when(backend.isStopping).thenReturn(true) | ||
| val allocator = newAllocator(backend, mock(classOf[ArmadaClient])) | ||
| primeDemand(allocator) | ||
|
|
||
| allocator.tryAllocateExecutors() | ||
|
|
||
| verify(backend, never()).getExecutorCounts | ||
| } | ||
|
|
||
| // Positive control: with the same primed demand but isStopping=false, the allocator must reach | ||
| // backend.getExecutorCounts. Proves the previous test isn't passing vacuously. | ||
| test("tryAllocateExecutors reaches getExecutorCounts when backend is not stopping") { | ||
| val backend = mock(classOf[ArmadaClusterManagerBackend]) | ||
| when(backend.isStopping).thenReturn(false) | ||
| when(backend.getExecutorCounts).thenReturn((4, 0)) | ||
| val allocator = newAllocator(backend, mock(classOf[ArmadaClient])) | ||
| primeDemand(allocator) | ||
|
|
There was a problem hiding this comment.
computeToAllocate property tests dropped without replacement
The entire table-driven test ("computeToAllocate caps the batch by gang cardinality") that shipped 10 parameterised cases — including the specific (4, 1, 1) case that was the original bug target — has been removed. The computeToAllocate logic is unchanged in ArmadaExecutorAllocator.scala, so it will regress silently if the gang-cap arithmetic is ever touched. Consider restoring the TableDrivenPropertyChecks-backed test alongside the new shutdown-gate tests.
Note: If this suggestion doesn't match your team's coding style, reply to this and let me know. I'll remember it for next time!
What: Restructure backend shutdown to reduce unnecessary executor cancellations and eliminate race-condition noise.
Why: During
stop(), the previous implementation calledsuper.stop()first (tearing down the RpcEnv), then cancelled all executor jobs. Executors that were still registering during the grace period hitRpcEndpointNotFoundExceptionand appeared as Failed in Armada Lookout. Additionally, bothonDisconnectedand the Armada event watcher could race to callremoveExecutorfor the same executor, producing noisy "Lost an executor (already removed)" ERROR logs.Changes:
stoppingflag withisStoppingaccessor toArmadaClusterManagerBackend; gatetryAllocateExecutorsandsubmitExecutorJobsinArmadaExecutorAllocatorso no new submissions happen during shutdownstop()order: stop allocator first so as to not be submitting any new pods while stopping, sendStopExecutorRPCs (keeping RpcEnv alive for late registrations), wait grace period, cancel remaining jobs, then tear down parent/event watcherRegisterExecutorduring shutdown: acknowledge success and immediately reply withStopExecutorso late executors exit 0 (Succeeded in Lookout) instead of crashing (Failed)Tests:
isStoppinglifecycle, pre-mark value correctness,safeRemoveExecutordeduplication, and allocator shutdown gates.Thanks
Thanks to @RammBagg for starting this PR.