Skip to content

Restructure backend shutdown to reduce unhelpful churn - #150

Open
GeorgeJahad wants to merge 3 commits into
masterfrom
fixCleanRace2
Open

GeorgeJahad wants to merge 3 commits into
masterfrom
fixCleanRace2

Conversation

@GeorgeJahad

@GeorgeJahad GeorgeJahad commented Jul 2, 2026 •

Copy link
Copy Markdown
Collaborator

What: Restructure backend shutdown to reduce unnecessary executor cancellations and eliminate race-condition noise.

Why: During stop(), the previous implementation called super.stop() first (tearing down the RpcEnv), then cancelled all executor jobs. Executors that were still registering during the grace period hit RpcEndpointNotFoundException and appeared as Failed in Armada Lookout. Additionally, both onDisconnected and the Armada event watcher could race to call removeExecutor for the same executor, producing noisy "Lost an executor (already removed)" ERROR logs.

Changes:

  • Add stopping flag with isStopping accessor to ArmadaClusterManagerBackend; gate tryAllocateExecutors and submitExecutorJobs in ArmadaExecutorAllocator so no new submissions happen during shutdown
  • Restructure stop() order: stop allocator first so as to not be submitting any new pods while stopping, send StopExecutor RPCs (keeping RpcEnv alive for late registrations), wait grace period, cancel remaining jobs, then tear down parent/event watcher
  • Intercept RegisterExecutor during shutdown: acknowledge success and immediately reply with StopExecutor so late executors exit 0 (Succeeded in Lookout) instead of crashing (Failed)

Tests:

  • Added tests for isStopping lifecycle, pre-mark value correctness, safeRemoveExecutor deduplication, and allocator shutdown gates.

Thanks
Thanks to @RammBagg for starting this PR.

RammBagg and others added 3 commits July 2, 2026 05:32
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>
@datadog-armadaproject

Copy link
Copy Markdown

Pipelines

⚠️ Warnings

🚦 1 Pipeline job failed

CI | E2E / End-to-end test 3.3.4 Scala 2.13.8   View in Datadog   GitHub Actions

This comment will be updated automatically if new data arrives.
🔗 Commit SHA: 2a82b6c | Docs | Give us feedback!

@greptile-apps

greptile-apps Bot commented Jul 2, 2026 •

Copy link
Copy Markdown
Contributor

Greptile Summary

This PR restructures ArmadaClusterManagerBackend.stop() to eliminate two classes of shutdown noise: spurious RpcEndpointNotFoundException failures from executors that register while the driver is shutting down, and duplicate "Lost an executor (already removed)" errors from onDisconnected racing the Armada event watcher.

  • Adds a @volatile stopping flag gated in the RegisterExecutor interceptor and the allocator's allocation/submission loops, so new executor submissions and full registrations are blocked as soon as shutdown begins. The stop() order is restructured to keep the RpcEnv alive through the grace period — stopExecutors() fires first, late RegisterExecutor messages receive an immediate StopExecutor reply (exiting 0), and super.stop() is deferred until after Armada jobs are cancelled.
  • Adds a removedExecutors ConcurrentHashMap key set to make safeRemoveExecutor idempotent, preventing duplicate removeExecutor calls when both onDisconnected and the Armada event watcher report the same executor exit. executorsPendingToRemove is now pre-populated with true across all kill/succeed/cancel paths so CoarseGrainedSchedulerBackend rewrites loss reasons to ExecutorKilled instead of counting clean exits as unexpected.

Confidence Score: 4/5

The 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

Filename Overview
src/main/scala/org/apache/spark/deploy/armada/Config.scala Grace period default bumped from 5s to 10s to allow more time for executor graceful exit during shutdown.
src/main/scala/org/apache/spark/scheduler/cluster/armada/ArmadaClusterManagerBackend.scala Core shutdown restructuring: adds stopping flag, removedExecutors dedup set, reorders stop() steps to keep RpcEnv alive during grace period, intercepts late RegisterExecutor during shutdown, and refines onDisconnected to handle pre-marked executors correctly.
src/main/scala/org/apache/spark/scheduler/cluster/armada/ArmadaExecutorAllocator.scala Adds isStopping early-return gates to tryAllocateExecutors and submitExecutorJobs; widens visibility to private[armada] for testability.
src/test/scala/org/apache/spark/scheduler/cluster/armada/ArmadaClusterManagerBackendSuite.scala New tests for isStopping lifecycle, executorsPendingToRemove value correctness, and safeRemoveExecutor deduplication; also adds short grace period to the before{} block to keep the suite fast.
src/test/scala/org/apache/spark/scheduler/cluster/armada/ArmadaExecutorAllocatorSuite.scala Replaces the table-driven computeToAllocate property test with new isStopping gate tests; the existing computeToAllocate coverage is dropped entirely.

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()
Loading
%%{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()
Loading

Comments Outside Diff (1)

  1. src/main/scala/org/apache/spark/scheduler/cluster/armada/ArmadaClusterManagerBackend.scala, line 555-566 (link)

    P2 removedExecutors gate marks on exception, preventing any retry

    removedExecutors.add(executorId) is evaluated before removeExecutor is called, so if removeExecutor throws a NonFatal exception the catch block silences it but the executor is permanently marked as removed. Any subsequent path that reaches safeRemoveExecutor for the same ID (e.g., onDisconnected firing just after a failed onExecutorSucceeded callback) will see a no-op. During normal shutdown the endpoint disappears cleanly so this is acceptable, but if removeExecutor fails mid-run (transient RPC error, not a shutdown-induced NPE), the executor will silently leak in Spark's internal state with no further attempt to clean it up.

Reviews (1): Last reviewed commit: "fix config" | Re-trigger Greptile

Comment on lines 22 to +70
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)

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.

P2 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!

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants