Skip to content

[SPARK-59542] Support spec.suspend to hold driver and master/worker creation in AppInitStep and ClusterInitStep - #828

Closed
dongjoon-hyun wants to merge 5 commits into
apache:mainfrom
dongjoon-hyun:SPARK-59542
Closed

dongjoon-hyun wants to merge 5 commits into
apache:mainfrom
dongjoon-hyun:SPARK-59542

Conversation

@dongjoon-hyun

@dongjoon-hyun dongjoon-hyun commented Sep 15, 2026 •

Copy link
Copy Markdown
Member

What changes were proposed in this pull request?

This PR aims to honor spec.suspend in AppInitStep and ClusterInitStep.

Resource State spec.suspend: true
SparkApplication Submitted, ScheduledToRestart Skip driver creation, completeAndDefaultRequeue()
SparkCluster Submitted Skip master / worker creation, completeAndDefaultRequeue()
  • The status is not written while suspended. A resource created with suspend: true therefore has no .status on the API server until it resumes, while an application held in ScheduledToRestart keeps the status its previous attempt already wrote. A dedicated Suspended state is planned in a follow-up.
  • The hold applies only when the driver pod (or the master StatefulSet) of the current attempt does not exist yet. If the resources were created but the status update failed, the initialization completes even if suspend is set meanwhile, so a live driver is never left unobserved. The current attempt's driver is the driver-labeled pod that has the name of the desired driver pod spec and is not terminating. The clean-up step deletes the previous attempt's driver before the application is scheduled to restart, so a previous attempt's pod is either gone or terminating even when the pod name is reused across attempts (e.g. user-specified spark.app.id). Because the informer cache may still hold the pre-deletion snapshot of that pod, the candidate is verified against the API server and the hold is bypassed only if the pod exists there and is not terminating; an API error keeps the hold for this reconciliation.
  • Setting spec.suspend back to false triggers a reconcile and enters the existing init path.
  • Validation, cleanup and deletion run before the init steps and are unaffected by spec.suspend.
  • Setting spec.suspend: true on a running application does not stop the current attempt. If the application is configured to restart, the next attempt is held in ScheduledToRestart. A running cluster is not affected.
  • The attempt duration used by restartCounterResetMillis is now measured from the first state after Submitted / ScheduledToRestart (normally DriverRequested), so time spent suspended or in restart backoff no longer counts as a successful run.
  • Add tests/e2e/suspend and a Suspend section to docs/spark_custom_resources.md.
  • Parameterize the resource name of the shared assertions in tests/e2e/assertions/ so that tests/e2e/suspend can reuse them.

Suspending a running attempt itself is out of scope and will be handled in a follow-up.

Why are the changes needed?

spec.suspend was added in SPARK-59475 but the operator does not read it yet. This is a prerequisite for the Kueue integration: job queueing systems such as Kueue create a workload with suspend: true and flip it to false once quota is available. The operator does not wire KueueWorkloadFactory into workload creation or admission handling yet.

Does this PR introduce any user-facing change?

Yes, in one place. spec.suspend itself is not released yet, but the attempt duration used by restartCounterResetMillis (released in 1.0.0) is now measured from DriverRequested instead of Submitted / ScheduledToRestart. Time spent in restart backoff (and suspended) no longer counts as a successful run, so an attempt whose running time is within restartBackoffMillis of restartCounterResetMillis no longer resets the restart counters and may reach maxRestartAttempts earlier than before. This only affects applications with restartCounterResetMillis >= 0 (default -1).

spec.suspend Before After
true Ignored, resources created immediately Held, no resources created
false Resources created Resources created (unchanged)
restartCounterResetMillis attempt duration Before After
Measured from Submitted / ScheduledToRestart (includes backoff) DriverRequested (excludes backoff and suspended time)

How was this patch tested?

Pass the CIs.

Test Coverage
AppInitStepTest Hold in Submitted and ScheduledToRestart, resume after suspend: false, proceed in RunningHealthy, recover when the driver was created but the status update failed, previous attempt's driver pod does not bypass the hold, stale informer snapshot of a deleted pod does not bypass the hold
SparkAppContextTest Current attempt's driver is found behind a previous attempt's pod; a terminating pod with the same name, a stale pre-deletion snapshot absent from the API server, a pod terminating on the API server, an API error during verification, or a pod with another name is not the current attempt's driver
ClusterInitStepTest Hold in Submitted, proceed in RunningHealthy, recover when the master StatefulSet already exists
ApplicationStatusTest Time held in ScheduledToRestart does not reset the restart counter (trimmed and untrimmed history)
tests/e2e/suspend Suspended app / cluster has no .status, no pod / StatefulSet, then patch and run to completion

Was this patch authored or co-authored using generative AI tooling?

Generated-by: Claude Fable 5.1

…reation in `AppInitStep` and `ClusterInitStep`

@peter-toth peter-toth left a comment •

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.

Thanks for the PR, @dongjoon-hyun!

AppInitStep and ClusterInitStep now return completeAndDefaultRequeue() before requesting the driver pod or the master / worker StatefulSets. The gate sits in the right place. Validation, cleanup and deletion all run ahead of it, and the restart backoff is skipped rather than reset. Two things block for me. Nothing in a suspended reconcile writes .status, so the two new *-suspended.yaml assertions have no status to match and the e2e group fails. And the doc's "no effect on a running application" does not hold for an app with a restart policy.

Blocking

  • 1. Suspended resource never gets a .status: every step of the suspended pipeline is a no-op on the status recorder, and toUpdateControl returns noUpdate(), so the API server keeps the resource with no .status at all. tests/e2e/suspend asserts status.currentState and status.stateTransitionHistory for both the app and the cluster, so the new group fails. [inline: tests/e2e/suspend/spark-application-suspended.yaml:25]
  • 2. Doc contradicts the ScheduledToRestart hold: "Setting it to true on a running application or cluster has no effect" is wrong for an app configured to restart. The current attempt keeps running, but the next one is held indefinitely, which is what suspendedAppScheduledToRestartDoesNotRequestDriver asserts. [inline: docs/spark_custom_resources.md:549]

Non-blocking

  • 3. No app-side test for the non-initializing case: ClusterInitStepTest.nonInitializingClusterProceeds pins that a RunningHealthy cluster ignores suspend. AppInitStepTest has no twin, and finding 2 is about exactly that boundary. [inline: AppInitStepTest.java:444]

Minor

namespace: default
spec:
suspend: true
status:

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.

Finding 1. Nothing in a suspended reconcile ever writes .status to the API server, so this assertion (and the spark-cluster-suspended.yaml twin) has no status to match.

The trace, for a SparkApplication created with suspend: true:

  • AppValidateStep persists only when the status is invalid, and it never is. CustomResource's constructor calls initStatus(), so getStatus() returns a Submitted ApplicationStatus even for a CR the server stored without one. isValidApplicationStatus passes and the step returns proceed().
  • AppCleanUpStep returns proceed() for Submitted without persisting.
  • the new branch returns completeAndDefaultRequeue() before any persist.
  • ReconcilerUtils.toUpdateControl returns UpdateControl.noUpdate(), and the CRD declares subresources: status: {}, so JOSDK writes nothing either.

SparkCluster is the same, except ClusterValidateStep is an unconditional proceed(), so there is not even a reset path.

I ran the three steps of each pipeline against a mock recorder:

SparkApplication app = new SparkApplication();
app.setMetadata(new ObjectMetaBuilder().withName("a").withNamespace("default").build());
app.getSpec().setSuspend(true);
SparkAppContext ctx = mock(SparkAppContext.class);
when(ctx.getResource()).thenReturn(app);
SparkAppStatusRecorder recorder = mock(SparkAppStatusRecorder.class);

new AppValidateStep().reconcile(ctx, recorder);
new AppCleanUpStep().reconcile(ctx, recorder);
new AppInitStep().reconcile(ctx, recorder);

verifyNoInteractions(recorder);  // passes; the ClusterValidate/Terminated/Init trio passes too

So a held resource has no status on the server and an empty Current State printer column. Chainsaw then fails here on status.currentState, and on (*.currentStateSummary) over a missing stateTransitionHistory. I have no cluster here, so that last hop is read from the assertion files rather than observed.

Two ways out:

  • Persist the initial status once when entering the hold, so the resource really is Submitted on the server. Careful: StatusRecorder.updateStatusFromCache seeds statusCache with the current status on a cache miss, and patchAndStatusWithVersionLocked short-circuits on newStatusNode.equals(previousStatusNode). A plain persistStatus(context, app.getStatus()) in the new branch is therefore dropped silently on the first reconcile.
  • Or drop the status: blocks from spark-application-suspended.yaml and spark-cluster-suspended.yaml and assert only spec.suspend: true plus the existing error: checks on the driver pod and the StatefulSets.

The first also fixes the operator-facing half. With Kueue holding a queue of workloads, kubectl get sparkapp currently shows a blank state for every one of them.

Comment thread docs/spark_custom_resources.md Outdated
```

* `suspend` only takes effect before the driver (or master / worker) resources are requested.
Setting it to `true` on a running application or cluster has no effect in the current version.

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.

Finding 2. Not true for a SparkApplication configured with a restart policy.

AppCleanUpStep has no suspend check, so a running attempt that fails still goes through terminateOrRestart and lands in ScheduledToRestart (AppCleanUpStep.java:186-193). The next reconcile hits the new gate and holds there until someone flips the flag back. suspendedAppScheduledToRestartDoesNotRequestDriver asserts exactly that hold. So suspend: true on a running app does not stop the current attempt, but it does stop every attempt after it.

For SparkCluster the bullet is accurate: SparkClusterReconciler.getReconcileSteps only adds ClusterInitStep for Submitted, and a RunningHealthy cluster never goes back there.

Suggested change
Setting it to `true` on a running application or cluster has no effect in the current version.
Setting it to `true` on a running application does not stop the current attempt. If the
application is configured to restart, the next attempt is held until `suspend` is set back to
`false`. Setting it to `true` on a running cluster has no effect in the current version.

The PR description's "Suspending an already running application or cluster is out of scope" reads the same way and is worth the same correction.

Assertions.assertEquals(
ApplicationStateSummary.DriverRequested,
application.getStatus().getCurrentState().getCurrentStateSummary());
}

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.

Finding 3. Worth an app-side twin of ClusterInitStepTest.nonInitializingClusterProceeds: a RunningHealthy app with suspend: true must still return proceed(). That is the boundary finding 2 is about, and right now it is pinned only for SparkCluster.

I ran this against the PR head and it passes:

@Test
void suspendedNonInitializingAppProceeds() {
  AppInitStep appInitStep = new AppInitStep();
  SparkAppContext mockContext = mock(SparkAppContext.class);
  SparkAppStatusRecorder recorder = mock(SparkAppStatusRecorder.class);
  SparkApplication application = new SparkApplication();
  application.setMetadata(applicationMetadata);
  application.getSpec().setSuspend(true);
  application.setStatus(
      application
          .getStatus()
          .appendNewState(
              new ApplicationState(ApplicationStateSummary.RunningHealthy, "running")));
  when(mockContext.getResource()).thenReturn(application);

  Assertions.assertEquals(
      ReconcileProgress.proceed(), appInitStep.reconcile(mockContext, recorder));
  verifyNoInteractions(recorder);
}

}
SparkCluster cluster = context.getResource();
if (cluster.getSpec().isSuspend()) {
log.info("Cluster is suspended, master and worker resources would not be requested.");

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.

Finding 4. A held resource logs this on every reconcile, so once per spark.kubernetes.operator.reconciler.intervalSeconds (120s by default) for as long as it stays suspended. With Kueue holding a few hundred workloads that is steady INFO traffic for a no-op. The other steady-state no-op paths report at debug (AppCleanUpStep.java:132, AppReconcileStep.java:86). Same line at AppInitStep.java:71.

spec:
suspend: false
status:
stateTransitionHistory:

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.

Finding 5. This is tests/e2e/assertions/spark-application/spark-state-transition.yaml with a different resource name, and spark-cluster-state-transition.yaml is the cluster one. AGENTS.md puts shared assertions in tests/e2e/assertions/. Binding the name there (name: ($SPARK_APPLICATION_NAME), as the namespace already is) would let this group reuse both files the way state-transition/chainsaw-test.yaml does, leaving only the two *-suspended.yaml files here. The extra spec.suspend: false check can move into an inline resource: assert.

@dongjoon-hyun

Copy link
Copy Markdown
Member Author

Thank you for the review, @peter-toth. All five findings are addressed in 5351a9a.

  1. No .status while suspended: confirmed. updateStatusFromCache seeds the cache with the in-memory Submitted status and patchAndStatusWithVersionLocked short-circuits on equality, so a plain persistStatus in the hold branch is dropped. I took the second option for this PR: the two *-suspended.yaml assertions now check spec.suspend: true plus the existing error: checks on the driver pod and the StatefulSets. A dedicated Suspended state, which also fixes the blank state column, is planned as a follow-up together with the running-attempt suspend.
  2. Doc: applied your suggestion. The PR description is updated the same way.
  3. App-side test: added AppInitStepTest.suspendedNonInitializingAppProceeds.
  4. Log level: both hold branches now log at debug.
  5. Shared assertions: spark-state-transition.yaml, spark-cluster-state-transition.yaml and spark-cluster-worker-statefulset.yaml now take ($SPARK_APPLICATION_NAME) / ($SPARK_CLUSTER_NAME), and tests/e2e/suspend reuses them. The watched-namespaces and watched-namespaces-file groups pass the new binding explicitly. The spec.suspend: false check moved into an inline resource: assert.

@dongjoon-hyun

Copy link
Copy Markdown
Member Author

Could you review this PR when you have some time, @viirya ?

@peter-toth peter-toth left a comment •

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.

Re-checked through 5351a9a. Findings 1, 2, 3, 4 and 5 are resolved and nothing regressed.

I re-ran the six new unit tests against ac5b0ee with both init steps reverted. Four fail there. The two nonInitializing*Proceeds guards pass on base, which is what a guard should do.

tests/e2e/state-transition already bound SPARK_APPLICATION_NAME and SPARK_CLUSTER_NAME at scenario level, so the parameterized shared assertions needed no change there. watched-namespaces and watched-namespaces-file were the only two consumers that did.

Blocking

  • 6. Docs promise a Submitted state the user cannot see (late catch): the new section says the operator "keeps the resource in its initializing state (Submitted, ...)". A resource created with suspend: true has no .status on the API server at all. kubectl get sparkapp shows an empty Current State, which is this repo's own signature for "not managed". One sentence in the docs fixes it. inline: docs/spark_custom_resources.md:533

Non-blocking

  • 7. The two *-suspended.yaml assertions cannot fail (new): each file's entire body is spec.suspend: true. That is the value the step applied two lines earlier, on a field the operator never writes. The error: checks carry the step on their own. inline: tests/e2e/suspend/spark-application-suspended.yaml:24
  • 8. The hold is keyed on state, not on whether the driver was already requested (late catch): if a reconcile creates the driver pod and then fails to persist DriverRequested, the app is back at Submitted with a live driver. Flipping suspend on at that point holds it there forever, and no observer runs in Submitted. Probably best carried into the running-attempt suspend you have planned. inline: .../AppInitStep.java:70

Minor

  • 9. Kueue bullet describes an integration that is not wired yet (late catch): KueueWorkloadFactory has no caller in src/main, so nothing creates a Workload and nothing flips suspend back. inline: docs/spark_custom_resources.md:554

Both `SparkApplication` and `SparkCluster` support `.spec.suspend`. When it is set to `true`, the
operator keeps the resource in its initializing state (`Submitted`, or `ScheduledToRestart` for an
application that is scheduled to restart) and does not request the driver pod or the master / worker
StatefulSets. Setting it back to `false` resumes the regular lifecycle.

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.

Finding 6. Submitted is not something the user can observe here. With the dedicated Suspended state deferred to a follow-up (#issuecomment-5684180823), the docs are the only place that can say so.

For a SparkApplication created with suspend: true, no step in the pipeline writes a status. AppValidateStep persists only when the status is invalid and it never is. AppCleanUpStep returns proceed() for Submitted. The new branch returns before any persist. toUpdateControl is noUpdate(). So .status is absent on the API server and the Current State printer column is blank.

That is indistinguishable from "the operator is not watching this namespace". tests/e2e/watched-namespaces/chainsaw-test.yaml:54-58 asserts .status == null as the signature of exactly that case.

Worth stating precisely, because the paragraph is only wrong for one of the two holds. An application suspended later, in ScheduledToRestart, does have a status on the server. AppCleanUpStep wrote it before the app got there.

Suggested change
StatefulSets. Setting it back to `false` resumes the regular lifecycle.
StatefulSets. Setting it back to `false` resumes the regular lifecycle.
`Submitted` here is the operator's in-memory view. A resource created with `suspend: true` has no
`.status` on the API server at all, so `kubectl get` shows an empty `Current State` until it
resumes. An application suspended later, in `ScheduledToRestart`, keeps the status its previous
attempt already wrote.

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

This distinction is still missing from the documentation at the current head. A valid resource created with suspend: true does not have its initial status persisted, so Submitted is an in-memory state and the Current State column remains blank. An application held in ScheduledToRestart, however, retains the status already written by cleanup.

Deferring a dedicated Suspended state seems reasonable for this PR, but please document the observable behavior explicitly. For example:

For a valid resource created with suspend: true, the initial Submitted status is not persisted to the API server, so kubectl get shows an empty Current State until initialization resumes. An application held in ScheduledToRestart retains its previously persisted status.

The PR description's statement that a held resource has no .status should likewise be limited to the initial-submission case.

name: spark-job-suspend-test
namespace: default
spec:
suspend: true

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.

Finding 7. This file, and spark-cluster-suspended.yaml, now assert only spec.suspend: true. That is the value spark-example-suspend.yaml applied two steps earlier, on a field the operator never writes. The assertion cannot fail.

Dropping the status: blocks was the right call. What is left is a 24-line file that adds nothing on top of the error: checks below it.

Two options. Delete both files and let the error: checks carry the step. Or make them assert the hold, using the idiom already in the tree at tests/e2e/watched-namespaces/chainsaw-test.yaml:54-58:

    - script:
        timeout: 30s
        content:
          kubectl get sparkapplication spark-job-suspend-test -n default -o json | jq ".status"
        check:
          (contains($stdout, 'null')): true

That pins today's behavior positively rather than by the absence of a pod. It will need updating when the Suspended state lands, which seems right. That follow-up is exactly the kind of change this test should notice.

return proceed();
}
SparkApplication app = context.getResource();
if (app.getSpec().isSuspend()) {

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.

Finding 8. The gate asks "is the app initializing and suspended?", not "has anything been requested yet?". Those two differ in one window, and appInitStepShouldBeIdempotentWhenStatusUpdateFails (AppInitStepTest.java:206) exists because the codebase already knows about it.

The sequence that produces it:

  1. Reconcile N creates the pre-resources, the driver pod and the driver resources.
  2. attemptStatusUpdate fails to persist DriverRequested and returns completeAndImmediateRequeue().
  3. StatusRecorder.persistStatus never updated statusCache, so updateStatusFromCache puts the resource back to Submitted on reconcile N+1.
  4. A driver pod is live and the app reads as Submitted.

An operator restart between the pod create and the status patch lands in the same place. The server has no .status, so initStatus() hands back a fresh Submitted.

If suspend goes to true in that window, this branch fires and the app never leaves Submitted. SparkAppReconciler.getReconcileSteps adds only AppInitStep for Submitted, so no observer and no timeout ever looks at that pod. It runs to completion unnoticed and the resource is held forever. That contradicts docs/spark_custom_resources.md:548, "suspend only takes effect before the driver ... resources are requested".

I reproduced it as a unit test at head. Reconcile 1 with persistStatus returning false, reset the status the way updateStatusFromCache would, set suspend, reconcile 2:

ReconcileProgress p2 = appInitStep.reconcile(ctx, recorder);
Assertions.assertEquals(ReconcileProgress.completeAndDefaultRequeue(), p2);
Assertions.assertEquals(
    ApplicationStateSummary.Submitted,
    application.getStatus().getCurrentState().getCurrentStateSummary());
// the driver pod from reconcile 1 is still there
Assertions.assertNotNull(
    kubernetesClient.pods().inNamespace("default").withName("driver-pod").get());

It passes, so the hold wins over the live driver.

The narrow fix here is if (app.getSpec().isSuspend() && context.getDriverPod().isEmpty()), which makes the code match the documented contract. One caveat I could not rule out. getDriverPod() reads the informer cache, and in ScheduledToRestart a just-deleted pod from the previous attempt could still be in it, which would make the operator start a driver for a suspended app.

If that risk is not worth taking in this PR, the running-attempt suspend you have planned is the natural home for this window. A suspend that can stop a live attempt stops this one too. Worth carrying into that work rather than leaving it only in this thread.

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

I agree with this finding, but I think it should be addressed before merging rather than deferred to support for suspending running attempts.

Resource creation and status persistence are separate operations. If the driver is created but persisting DriverRequested fails, the next reconciliation restores Submitted from the status cache. An operator restart between those operations produces the same discrepancy.

If suspend becomes true at that point, this branch prevents initialization recovery. SparkAppReconciler selects only validation, cleanup, and initialization for Submitted, so the existing driver is no longer observed for completion or timeout while the flag remains set. The equivalent window also exists in ClusterInitStep between resource creation and persisting RunningHealthy.

Could we distinguish an initialization that has already started from one that has not requested resources yet, and allow the former to recover? Please add a regression test covering successful resource creation, failed status persistence, and then enabling suspend.

I would avoid fixing this solely with getDriverPod().isEmpty(): that lookup uses the informer cache, and the driver labels do not distinguish attempts, so a previous attempt's pod could incorrectly bypass the hold.

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

The assumption that the driver pod name carries the attempt ID does not hold for all supported configurations.

SparkAppSubmissionWorker.buildDriverConf() preserves a user-specified spark.app.id through setIfMissing(), and SparkAppDriverConf.resourceNamePrefix() returns that ID. The existing checkAppIdWhenUserSpecifiedInSparkConf test explicitly covers this behavior. Spark also allows spark.kubernetes.driver.pod.name to override the pod name directly.

Consequently, consecutive attempts can have the same expected pod name. If the previous attempt's pod remains in the informer cache after cleanup, this comparison returns true for a suspended ScheduledToRestart application. Once backoff has elapsed, initialization proceeds; if the old pod has already disappeared from the API server, it can create a new driver despite suspend: true.

Could we identify the current attempt independently of a potentially reused pod name, or explicitly validate any naming restrictions required by this approach? Please extend previousAttemptDriverPodDoesNotBypassSuspend to cover a fixed spark.app.id and an explicitly configured driver pod name. Its current use of two different names does not exercise this case.

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

This handles a previous attempt's pod once the informer has observed its deletion, but the stale pre-deletion snapshot can still pass this filter.

The remaining sequence is:

  1. Consecutive attempts reuse the driver pod name.
  2. Cleanup successfully deletes the old pod and persists ScheduledToRestart.
  3. With suspend: true and zero or elapsed backoff, the next reconcile still sees the old informer snapshot: the expected name and a null deletionTimestamp.
  4. This helper returns the old pod, so AppInitStep bypasses the hold.
  5. Initialization queries the API server, finds that the pod is gone, and creates a new driver despite suspend.

I agree that this is informer lag, but the missing-driver grace period does not make this path safe. That flow defers a decision and eventually performs live verification; this path immediately authorizes resource creation. The grace-period observer also does not run for ScheduledToRestart.

Could we perform a live lookup by the desired pod name before allowing the suspend bypass, require the pod to exist and not be terminating, and requeue if verification fails?

The regression test should keep the pre-deletion pod snapshot in the informer while the API server reports the pod absent, then verify that a suspended application creates no driver. The current test sets deletionTimestamp on the cached object, so it covers the already-updated cache rather than this remaining window.

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

Confirmed resolved in 18636ac. The live GET prevents the stale pre-deletion snapshot from authorizing initialization, and lookup failures keep the hold. staleInformerSnapshotDoesNotBypassSuspend now covers the specific sequence discussed here and verifies that no driver is created and the application stays in ScheduledToRestart. Thanks for adding this coverage.

`false`. Setting it to `true` on a running cluster has no effect in the current version.
* Deleting a suspended resource works as usual.
* This is the building block for external job queueing systems such as
[Kueue](https://kueue.sigs.k8s.io/), which admit a workload by flipping `suspend` to `false`.

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.

Finding 9. Reads as if Kueue does this today. KueueWorkloadFactory is the only code outside the two init steps that touches spec.isSuspend(), and it has no caller in src/main. Nothing creates a Workload, and nothing flips suspend back.

Suggested change
[Kueue](https://kueue.sigs.k8s.io/), which admit a workload by flipping `suspend` to `false`.
[Kueue](https://kueue.sigs.k8s.io/), which admit a workload by flipping `suspend` to `false`.
The operator does not integrate with such a system yet.

@viirya

viirya commented Sep 15, 2026

Copy link
Copy Markdown
Member

I will take a look.

@viirya viirya left a comment •

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

Thanks for the update. I reviewed this against 5351a9a, including the earlier review threads and the interactions with status recovery, cleanup, and restart handling.

The two init gates are small and consistent with the existing reconciliation structure. The first-round findings have been addressed, and the CI run for this commit is green, including both suspend e2e jobs.

I still have two correctness concerns:

  1. If resource creation succeeds but the status update fails, enabling suspend prevents the existing initialization recovery path from completing. A driver can remain live while the application stays in Submitted, where no driver observers run.
  2. Time spent suspended is included in the attempt duration used by restartCounterResetMillis. A long hold can therefore reset retry counters even when the resumed attempt fails immediately.

I would address these interactions and add regression coverage before merging. The suspend documentation also needs to distinguish the in-memory Submitted state from the status actually visible on the API server.

Two smaller suggestions:

  • The spec.suspend: true assertions confirm the applied value, but do not establish that the operator processed the hold. The absence-of-resource checks and successful resume provide the behavioral coverage. Replacing these assertions with .status == null would not establish reconciliation either.
  • Please clarify that this is a prerequisite for future Kueue integration; the operator does not yet wire KueueWorkloadFactory into workload creation or admission handling.

Comment on lines +70 to +72
if (app.getSpec().isSuspend()) {
log.debug("Application is suspended, driver resources would not be requested.");
return completeAndDefaultRequeue();

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

Holding the application here also preserves the initializing state's timestamp. This affects the existing restart-counter logic: ApplicationStatus.calculateCurrentAttemptDuration() measures from the latest Submitted or ScheduledToRestart, and terminateOrRestart() applies the time-based counter reset before checking retry limits.

For example, with maxRestartAttempts: 1 and restartCounterResetMillis: 3600000, hold the permitted retry in ScheduledToRestart for over an hour, then resume it and let it fail immediately. The suspended hour satisfies the reset threshold, so the counter resets and another retry is allowed instead of exhausting the configured limit.

This makes time spent waiting for admission count toward the documented reset for a long-running attempt. Could we exclude suspended time from that calculation and cover this sequence with both trimmed and untrimmed histories?

Simply changing the initializing timestamp on resume would also change restart-backoff and history semantics, so those should be preserved when addressing this.

@dongjoon-hyun

Copy link
Copy Markdown
Member Author

Thank you for the second round, @peter-toth and @viirya. All findings are addressed in 34dbda1.

6. Docs / Submitted is in-memory (peter-toth, viirya): added a paragraph to the Suspend section. A resource created with suspend: true has no .status on the API server and kubectl get shows an empty Current State until it resumes, while an application held in ScheduledToRestart keeps the status its previous attempt wrote. The PR description is limited to the initial-submission case the same way.

7. *-suspended.yaml cannot fail (peter-toth, viirya): deleted both files and replaced them with the kubectl get ... | jq .status / contains($stdout, 'null') check already used in watched-namespaces. As @viirya noted, this pins today's observable behavior rather than proving the hold was reconciled; the behavioral coverage stays with the error: checks on the driver pod / StatefulSets and the resume to completion. The null check will need updating when the Suspended state lands, which seems right.

8. Hold keyed on state, not on whether the driver was requested (peter-toth, viirya): the hold now applies only when the driver pod of the current attempt does not exist. AppInitStep compares the informer pod's name with getDriverPodSpec()'s name rather than using getDriverPod().isEmpty(), so a pod left from a previous attempt (same labels, different attempt id in the name) does not bypass the hold. ClusterInitStep skips the hold when the master StatefulSet already exists. This also covers the operator-restart variant, since the check does not depend on the status cache. Regression tests: suspendAfterDriverRequestedCompletesInitialization (create succeeds, status persist fails, suspend enabled, next reconcile reaches DriverRequested), previousAttemptDriverPodDoesNotBypassSuspend, and suspendAfterMasterRequestedCompletesInitialization.

9. Kueue bullet (peter-toth, viirya): applied the suggestion. The PR description now says this is a prerequisite and that KueueWorkloadFactory is not wired into workload creation or admission yet.

Suspended time counted toward restartCounterResetMillis (viirya): calculateCurrentAttemptDuration() now measures from the first state after Submitted / ScheduledToRestart (normally DriverRequested) instead of the initializing state itself. The initializing timestamp and the history are untouched, so restart backoff keeps its base. Time spent suspended, and time spent in restart backoff, no longer counts as a successful run. Added testTimeHeldInScheduledToRestartDoesNotResetRestartCounter for both trimmed and untrimmed histories (2h hold in ScheduledToRestart, resume, immediate failure, maxRestartAttempts: 1 with a 1h reset window ends in ResourceReleased), and updated testCalculateCurrentAttemptDuration for the new start point. The Restart Counter reset docs mention the measurement point.

@viirya viirya left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

Thanks for addressing the previous comments. I re-reviewed 34dbda1.

The documentation now distinguishes the in-memory Submitted state from persisted status and clarifies the scope of the Kueue integration. The attempt-duration change also addresses the suspended-time issue while preserving the restart-backoff timestamp.

The initialization recovery change is an improvement, but I do not think Finding 8 is fully resolved yet. The new helper assumes that pod names uniquely identify attempts, which does not hold when users configure spark.app.id or spark.kubernetes.driver.pod.name. It also checks the name only after selecting an arbitrary driver pod from the informer cache.

Could we tighten the current-attempt lookup and extend the regression coverage for these cases before merging?

Comment on lines +165 to +167
private boolean isDriverRequested(SparkAppContext context) {
Optional<Pod> driverPod = context.getDriverPod();
return driverPod.isPresent()

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

context.getDriverPod() filters by application/driver labels and then calls findAny(). The name comparison here happens after that selection.

If both an older attempt's pod and the current attempt's pod are present in the informer cache, this can select the older one and return false without checking the current pod. In the status-persistence failure scenario this helper is intended to recover, that leaves the application held even though its current driver exists.

Could we apply the current-attempt predicate before selecting a pod, rather than checking only the single pod returned by getDriverPod()? A regression test with both pods available, with the older pod encountered first, would cover this. The existing tests mock only one returned pod, so they cannot detect this selection issue.

@dongjoon-hyun

Copy link
Copy Markdown
Member Author

Thank you for the careful follow-up, @viirya. Both points are addressed in 67d1cfc.

Selection before the name check: SparkAppContext now has getCurrentAttemptDriverPod(). It collects every driver-labeled pod from the informer cache and applies the current-attempt predicate to all of them before selecting one, instead of checking only the single pod that getDriverPod().findAny() happens to return. AppInitStep uses that method; the name comparison is gone from the step. SparkAppContextTest.currentAttemptDriverPodIsFoundBehindPreviousAttemptPod has the older pod first in the stream and asserts the current one is found.

Reused pod names (spark.app.id / spark.kubernetes.driver.pod.name): rather than assuming the name carries the attempt id, the predicate now requires the pod to have the desired driver pod spec's name and no deletionTimestamp. AppCleanUpStep deletes the previous attempt's driver before the application transitions to ScheduledToRestart, and deleteResourceIfExists rethrows anything but 404 so the transition cannot happen without the delete being accepted. Once the app is in ScheduledToRestart, the previous attempt's pod is therefore either gone or terminating, even when its name is identical to the current spec's. Submitted has no previous attempt. previousAttemptDriverPodDoesNotBypassSuspend now uses a terminating pod with the very same name as the spec, and SparkAppContextTest.terminatingPreviousAttemptPodWithSameNameIsNotCurrentAttemptDriver covers the selection for that case. The remaining window is an informer cache that has not yet reflected the accepted delete, which is the same class of lag the missing-driver grace period already tolerates.

I considered labeling the driver pod with the attempt id, which would make the identification exact regardless of naming. That changes every driver pod the operator creates and is independent of spec.suspend, so it is better done as its own change if we want it. The javadoc and the PR description describe the current rule.

The pi-with-gluten failure on the previous run was a state-transition timeout in that group's own assertion, unrelated to this change.

@viirya viirya left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

Thanks for the update. I re-reviewed 67d1cfc.

The selection-order issue is resolved: the current-attempt predicate is now applied before selecting a pod, and the new test covers the older pod appearing first. The documentation and restart-duration fixes also remain addressed.

I still think Finding 8 needs one more change. Filtering on the cached deletionTimestamp handles a terminating pod only after the informer has observed the deletion. It does not cover the stale pre-deletion snapshot described in the previous comment.

That distinction matters here because a false positive allows initialization to create a new driver despite suspend: true. Could we verify the pod against the API server before using its existence to bypass suspend, and add a regression test for that stale-snapshot case? This can stay local to the suspend recovery path without changing labels on every driver pod.

@dongjoon-hyun

Copy link
Copy Markdown
Member Author

Thank you, @viirya. You are right that the cached deletionTimestamp only covers a cache that has already observed the delete, and that this path authorizes creation immediately rather than deferring. Addressed in 18636ac.

Live verification before the bypass: SparkAppContext.getCurrentAttemptDriverPod() still uses the informer cache to find a candidate (driver labels, the desired driver pod name, no deletionTimestamp), but it now verifies that candidate against the API server with a GET by the desired pod name and returns it only if the pod exists there and is not terminating. A stale pre-deletion snapshot therefore no longer bypasses the hold: the live lookup finds nothing, AppInitStep keeps the hold and requeues at the regular interval. A KubernetesClientException during the lookup is logged and treated as absent, so an API error keeps the hold for that reconciliation instead of authorizing creation. When the cache has no candidate, neither the driver spec is built nor the API is called, so a held application still costs nothing per reconcile.

Regression coverage:

  • AppInitStepTest.staleInformerSnapshotDoesNotBypassSuspend runs the step against a real SparkAppContext whose informer stream returns the pre-deletion snapshot (same name, no deletionTimestamp) while the mock API server has no such pod, with suspend: true and an elapsed backoff in ScheduledToRestart. It asserts that no driver pod is created, no status is persisted and the app stays in ScheduledToRestart.
  • SparkAppContextTest adds stalePreDeletionSnapshotIsNotCurrentAttemptDriver, podTerminatingOnApiServerIsNotCurrentAttemptDriver and apiErrorDuringVerificationIsTreatedAsAbsent, and checks that no API call is made when the cache has no matching candidate.

The javadoc and the PR description describe the verification step. The Java 21 CI failure on the previous run was SparkOperatorConfigMapReconcilerTest.initializationError failing to download the kubeapitest binary, unrelated to this change.

@viirya viirya left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

Thanks for addressing the review feedback. I re-reviewed 18636ac, including the earlier discussion threads.

The live lookup now closes the stale-cache case: a cached pod only permits the suspend bypass after the API server confirms that the pod exists and is not terminating. Verification errors keep the application suspended and requeue.

The new regression test covers the pre-deletion informer snapshot with the pod already absent from the API server, using the real SparkAppContext. The selection-order, restart-duration, and documentation concerns are also addressed.

All 35 CI jobs for this commit passed, including the suspend e2e tests on Kubernetes 1.34 and 1.37.

I have no remaining blocking concerns. LGTM.

@dongjoon-hyun

Copy link
Copy Markdown
Member Author

Thank you, @viirya and @peter-toth ! This is the first step and I'm going to add future steps. It will become more mature than the AS-IS status.

@dongjoon-hyun dongjoon-hyun added this to the 1.1.0 milestone Sep 18, 2026
@dongjoon-hyun

Copy link
Copy Markdown
Member Author

Merged to main

@dongjoon-hyun
dongjoon-hyun deleted the SPARK-59542 branch September 18, 2026 18:26
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.

3 participants