Skip to content
Closed
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
5 changes: 5 additions & 0 deletions .github/workflows/build_and_test.yml
Original file line number Diff line number Diff line change
Expand Up @@ -91,6 +91,7 @@ jobs:
- resource-selector
- watched-namespaces
- driver-start-timeout
- suspend
exclude:
- mode: dynamic
test-group: spark-versions
Expand All @@ -106,6 +107,8 @@ jobs:
test-group: resource-selector
- mode: dynamic
test-group: driver-start-timeout
- mode: dynamic
test-group: suspend
- mode: static
test-group: watched-namespaces
- mode: static
Expand All @@ -124,6 +127,8 @@ jobs:
test-group: watched-namespaces
- mode: selector
test-group: driver-start-timeout
- mode: selector
test-group: suspend
include:
- kubernetes-version: "1.37.0"
mode: dynamic-file
Expand Down
38 changes: 37 additions & 1 deletion docs/spark_custom_resources.md
Original file line number Diff line number Diff line change
Expand Up @@ -380,7 +380,9 @@ restartConfig:
The `restartCounterResetMillis` field controls automatic restart counter resets for long-running
application attempts. When set to a non-negative value (in milliseconds), the operator will reset
all restart counters (including the general counter and both failure counters) if an application
attempt runs successfully for at least the specified duration before ending.
attempt runs successfully for at least the specified duration before ending. The duration is
measured from the first state after `Submitted` / `ScheduledToRestart` (normally
`DriverRequested`), so time spent in restart backoff or suspended is not counted.

Time-based reset takes highest precedence over all limit checks. If an attempt runs longer than
`restartCounterResetMillis`, the operator will always restart with reset counters, regardless
Expand Down Expand Up @@ -525,6 +527,40 @@ Note that `ttlAfterStopMillis` applies to the app as well as its secondary resou
latter is smaller, then it takes higher precedence: operator would remove all resources related
to this app after `ttlAfterStopMillis`.

## Suspend

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.


`Submitted` here is the operator's in-memory view. 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
later, in `ScheduledToRestart`, keeps the status its previous attempt already wrote.

``` yaml
apiVersion: spark.apache.org/v1
kind: SparkApplication
metadata:
name: suspended-pi
spec:
suspend: true
mainClass: "org.apache.spark.examples.SparkPi"
jars: "local:///opt/spark/examples/jars/spark-examples.jar"
runtimeVersions:
sparkVersion: "4.2.0"
```

* `suspend` only takes effect before the driver (or master / worker) resources are requested.
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.
* 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.

The operator does not integrate with such a system yet.

## Spark Cluster

Spark Operator also supports launching Spark clusters in k8s via `SparkCluster` custom resource,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -260,17 +260,26 @@ protected ApplicationState findFirstStateOfCurrentAttempt() {
/**
* Calculates the duration of the current application attempt.
*
* <p>The duration is calculated as the time between the first state of the current attempt (as
* determined by {@link #findFirstStateOfCurrentAttempt()}) and the current state's last
* transition time. This is particularly useful for determining whether the restart counter should
* be reset based on the configured {@code restartCounterResetMillis}.
* <p>The duration is calculated as the time between the first state after the initializing state
* of the current attempt (normally DriverRequested) and the current state's last transition time,
* so that time spent in restart backoff or suspended is not counted. If the history has no
* initializing state, the first entry is used. This is particularly useful for determining
* whether the restart counter should be reset based on the configured {@code
* restartCounterResetMillis}.
*
* @return A Duration representing the time elapsed since the start of the current attempt.
* @return A Duration representing the time elapsed since the current attempt started running.
*/
protected Duration calculateCurrentAttemptDuration() {
ApplicationState firstStateOfCurrentAttempt = findFirstStateOfCurrentAttempt();
List<ApplicationState> states = new ArrayList<>(stateTransitionHistory.values());
ApplicationState attemptStart = states.get(0);
for (int k = states.size() - 1; k >= 0; k--) {
if (states.get(k).getCurrentStateSummary().isInitializing()) {
attemptStart = states.get(Math.min(k + 1, states.size() - 1));
break;
}
}
return Duration.between(
Instant.parse(firstStateOfCurrentAttempt.getLastTransitionTime()),
Instant.parse(attemptStart.getLastTransitionTime()),
Instant.parse(currentState.getLastTransitionTime()));
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -354,15 +354,20 @@ void testCalculateCurrentAttemptDuration() throws Exception {
assertNotNull(durationMultipleStates);
assertTrue(durationMultipleStates.toMillis() >= 0);

// Test with restart scenario - duration should be calculated from ScheduledToRestart state
// Test with restart scenario - duration should be calculated from the first state after
// ScheduledToRestart (DriverRequested), excluding the time spent in backoff or suspended.
// Create states with explicit timestamps
Instant now = Instant.now();
Instant oneHourAgo = now.minus(Duration.ofHours(1));
Instant tenMinutesAgo = now.minus(Duration.ofMinutes(10));
Instant fiveMinutesAgo = now.minus(Duration.ofMinutes(5));

ApplicationState expectedSecondAttemptStart =
new ApplicationState(ApplicationStateSummary.ScheduledToRestart, "");
expectedSecondAttemptStart.setLastTransitionTime(tenMinutesAgo.toString());
ApplicationState secondAttemptDriverRequested =
new ApplicationState(ApplicationStateSummary.DriverRequested, "");
secondAttemptDriverRequested.setLastTransitionTime(fiveMinutesAgo.toString());
ApplicationState expectedSecondAttemptEnd =
new ApplicationState(ApplicationStateSummary.Failed, "");

Expand All @@ -372,23 +377,70 @@ void testCalculateCurrentAttemptDuration() throws Exception {
.appendNewState(new ApplicationState(ApplicationStateSummary.RunningHealthy, ""))
.appendNewState(new ApplicationState(ApplicationStateSummary.Failed, ""))
.appendNewState(expectedSecondAttemptStart)
.appendNewState(new ApplicationState(ApplicationStateSummary.DriverRequested, ""))
.appendNewState(secondAttemptDriverRequested)
.appendNewState(expectedSecondAttemptEnd);

// Verify it finds ScheduledToRestart as the first state of current attempt
ApplicationState firstState = statusWithRestarts.findFirstStateOfCurrentAttempt();
assertEquals(expectedSecondAttemptStart, firstState);

// Verify duration is calculated from ScheduledToRestart state (10 minutes ago)
// Verify duration is calculated from DriverRequested state (5 minutes ago)
Duration durationAfterRestart = statusWithRestarts.calculateCurrentAttemptDuration();
assertNotNull(durationAfterRestart);

Duration expectedDuration =
Duration.between(
tenMinutesAgo, Instant.parse(expectedSecondAttemptEnd.getLastTransitionTime()));
fiveMinutesAgo, Instant.parse(expectedSecondAttemptEnd.getLastTransitionTime()));
assertEquals(expectedDuration, durationAfterRestart);
}

@Test
void testTimeHeldInScheduledToRestartDoesNotResetRestartCounter() {
RestartConfig config = new RestartConfig();
config.setRestartPolicy(RestartPolicy.Always);
config.setMaxRestartAttempts(1L);
config.setRestartCounterResetMillis(3600000L); // 1 hour

Instant now = Instant.now();
Instant threeHoursAgo = now.minus(Duration.ofHours(3));
// The first attempt fails right away
ApplicationStatus status =
createInitialStatusWithSubmittedTime(threeHoursAgo)
.appendNewState(stateAt(ApplicationStateSummary.DriverRequested, threeHoursAgo))
.appendNewState(stateAt(ApplicationStateSummary.Failed, threeHoursAgo));

for (boolean trimStateTransitionHistory : new boolean[] {false, true}) {
ApplicationStatus restarted =
status.terminateOrRestart(
config, ResourceRetainPolicy.Never, null, trimStateTransitionHistory);
assertEquals(
ApplicationStateSummary.ScheduledToRestart,
restarted.getCurrentState().getCurrentStateSummary());
assertEquals(1L, restarted.getCurrentAttemptSummary().getAttemptInfo().getRestartCounter());

// The permitted retry is held (e.g. suspended) for two hours before it resumes and fails
// immediately: the time spent on hold must not count as a successful run.
restarted.getCurrentState().setLastTransitionTime(now.minus(Duration.ofHours(2)).toString());
ApplicationStatus resumed =
restarted
.appendNewState(stateAt(ApplicationStateSummary.DriverRequested, now))
.appendNewState(stateAt(ApplicationStateSummary.Failed, now));

ApplicationStatus terminated =
resumed.terminateOrRestart(
config, ResourceRetainPolicy.Never, null, trimStateTransitionHistory);
assertEquals(
ApplicationStateSummary.ResourceReleased,
terminated.getCurrentState().getCurrentStateSummary());
}
}

private ApplicationState stateAt(ApplicationStateSummary summary, Instant time) {
ApplicationState state = new ApplicationState(summary, "");
state.setLastTransitionTime(time.toString());
return state;
}

private ApplicationStatus createInitialStatusWithSubmittedTime(Instant submittedTime) {
ApplicationStatus status = new ApplicationStatus();
ApplicationState submittedState = status.getStateTransitionHistory().get(0L);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -30,6 +30,7 @@
import io.fabric8.kubernetes.api.model.HasMetadata;
import io.fabric8.kubernetes.api.model.Pod;
import io.fabric8.kubernetes.client.KubernetesClient;
import io.fabric8.kubernetes.client.KubernetesClientException;
import io.javaoperatorsdk.operator.api.reconciler.Context;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
Expand Down Expand Up @@ -73,6 +74,61 @@ public Optional<Pod> getDriverPod() {
.findAny();
}

/**
* Returns the driver pod of the current attempt, if present. A pod counts as the current
* attempt's driver only if it carries the driver labels, has the name of the desired driver pod
* spec and is not being deleted. The last condition tells a previous attempt's pod apart when
* the driver pod name is reused across attempts (e.g. user-specified spark.app.id): the previous
* attempt's pod has been deleted by the clean-up step before the application was scheduled to
* restart, so it is either gone or terminating.
*
* <p>The informer cache is only used to find a candidate. Since the cache may still hold the
* pre-deletion snapshot of a previous attempt's pod, the candidate is verified against the API
* server and returned only if it exists there and is not terminating. If the verification fails
* because of an API error, the pod is considered absent for this reconciliation.
*
* @return An Optional containing the driver Pod of the current attempt, or empty if not found.
*/
public Optional<Pod> getCurrentAttemptDriverPod() {
List<Pod> driverPods =
josdkContext
.getSecondaryResourcesAsStream(Pod.class)
.filter(
p ->
p.getMetadata()
.getLabels()
.entrySet()
.containsAll(driverLabels(sparkApplication).entrySet()))
.filter(p -> p.getMetadata().getDeletionTimestamp() == null)
.toList();
if (driverPods.isEmpty()) {
return Optional.empty();
}
String driverPodName = getDriverPodSpec().getMetadata().getName();
if (driverPods.stream().noneMatch(p -> driverPodName.equals(p.getMetadata().getName()))) {
return Optional.empty();
}
try {
Pod livePod =
josdkContext
.getClient()
.pods()
.inNamespace(sparkApplication.getMetadata().getNamespace())
.withName(driverPodName)
.get();
if (livePod == null || livePod.getMetadata().getDeletionTimestamp() != null) {
return Optional.empty();
}
return Optional.of(livePod);
} catch (KubernetesClientException e) {
log.warn(
"Failed to verify driver pod {} against the API server, considering it absent.",
driverPodName,
e);
return Optional.empty();
}
}

/**
* Live lookup of the driver pod against the API server, bypassing the informer cache. The result
* is memoized for the lifetime of this context so that multiple observe steps share a single
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -67,6 +67,10 @@ public ReconcileProgress reconcile(
return proceed();
}
SparkApplication app = context.getResource();
if (app.getSpec().isSuspend() && !isDriverRequested(context)) {
log.debug("Application is suspended, driver resources would not be requested.");
return completeAndDefaultRequeue();
}
if (app.getStatus().getPreviousAttemptSummary() != null) {
Instant lastTransitionTime = Instant.parse(currentState.getLastTransitionTime());
ApplicationAttemptSummary attemptSummary = app.getStatus().getPreviousAttemptSummary();
Expand Down Expand Up @@ -148,6 +152,20 @@ public ReconcileProgress reconcile(
return attemptStatusUpdate(context, statusRecorder, updatedStatus, completeAndDefaultRequeue());
}

/**
* Checks whether the driver pod of the current attempt has already been requested. This covers
* the case where the driver was created but the status update to DriverRequested failed, so that
* a suspended application still completes its initialization instead of being held with a live
* driver. See {@link SparkAppContext#getCurrentAttemptDriverPod()} for how a pod left from a
* previous attempt is told apart.
*
* @param context The SparkAppContext for the application.
* @return True if the driver pod of the current attempt exists, false otherwise.
*/
private boolean isDriverRequested(SparkAppContext context) {
return context.getCurrentAttemptDriverPod().isPresent();
}

/**
* Creates an ApplicationState indicating a scheduling failure due to resource creation failure.
*
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -39,6 +39,7 @@
import org.apache.spark.k8s.operator.reconciler.ReconcileProgress;
import org.apache.spark.k8s.operator.status.ClusterState;
import org.apache.spark.k8s.operator.status.ClusterStatus;
import org.apache.spark.k8s.operator.utils.ReconcilerUtils;
import org.apache.spark.k8s.operator.utils.SparkClusterStatusRecorder;

/** Request cluster master and its resources when starting an attempt. */
Expand All @@ -59,6 +60,14 @@ public ReconcileProgress reconcile(
return proceed();
}
SparkCluster cluster = context.getResource();
// A cluster whose master StatefulSet already exists has been requested before (e.g. the status
// update to RunningHealthy failed), so let it complete its initialization even if suspended.
if (cluster.getSpec().isSuspend()
&& ReconcilerUtils.getResource(context.getClient(), context.getMasterStatefulSetSpec())
.isEmpty()) {
log.debug("Cluster is suspended, master and worker resources would not be requested.");
return completeAndDefaultRequeue();
}
if (cluster.getStatus().getPreviousAttemptSummary() != null) {
Instant lastTransitionTime = Instant.parse(currentState.getLastTransitionTime());
Instant restartTime = lastTransitionTime.plusMillis(300 * 1000);
Expand Down
Loading