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
9 changes: 6 additions & 3 deletions docs/architecture.md
Original file line number Diff line number Diff line change
Expand Up @@ -138,8 +138,8 @@ stateDiagram-v2
RunningWithPartialCapacity --> Succeeded
RunningWithPartialCapacity --> Failed

RunningHealthy --> Suspended : spec.suspend=true
Suspended --> Submitted : spec.suspend=false, new attempt
RunningHealthy --> Suspended : spec.suspend=true or Kueue eviction
Suspended --> Submitted : spec.suspend=false or after Kueue eviction, new attempt

state Failures {
SchedulingFailure
Expand Down Expand Up @@ -171,7 +171,10 @@ stateDiagram-v2
* Likewise, an application moves to `Suspended` from any of the states from `DriverRequested` to
`RunningWithBelowThresholdExecutors` when [`spec.suspend`](spark_custom_resources.md#suspend) is
set to `true`. Its driver and executors are released, and setting it back to `false` starts a
new attempt from `Submitted`, which does not count as a restart.
new attempt from `Submitted`, which does not count as a restart. An application queued by
[Kueue](spark_custom_resources.md#kueue) is suspended as well when Kueue evicts its `Workload`,
and starts such a new attempt once its driver and executors are released and the requeue
backoff of its `Workload` elapsed, or once its `Workload` is reactivated if it was deactivated.
* User may configure the app CR to time-out after given threshold of time if it cannot reach healthy
state after given threshold. The timeout can be configured for different lifecycle stages,
when driver starting, when driver becoming ready, and when requesting executor pods.
Expand Down
2 changes: 1 addition & 1 deletion docs/config_properties.md
Original file line number Diff line number Diff line change
Expand Up @@ -53,7 +53,7 @@
| spark.kubernetes.operator.reconciler.retry.initialIntervalSeconds | Integer | 5 | false | Initial interval(in seconds) of retries on unhandled controller errors. |
| spark.kubernetes.operator.reconciler.retry.intervalMultiplier | Double | 1.5 | false | Interval multiplier of retries on unhandled controller errors. Setting this to 1 for linear retry. |
| spark.kubernetes.operator.reconciler.retry.maxIntervalSeconds | Integer | -1 | false | Max interval(in seconds) of retries on unhandled controller errors. Set to non-positive for unlimited. |
| spark.kubernetes.operator.reconciler.suspendHoldRequeueIntervalSeconds | Long | 1800 | true | Requeue interval (in seconds) at which a resource held by spec.suspend is reconciled, so that its SuspendHeld event is republished. A SparkApplication or SparkCluster suspended while running publishes no such event, since its Suspended status says so, and uses this interval only to look at its released resources again. The hold ends only when a user clears spec.suspend, which arrives as a watch event and reconciles right away, so nothing waits for this interval. It is deliberately much coarser than spark.kubernetes.operator.reconciler.intervalSeconds, since a suspended resource has nothing to observe while each republish costs a read and a write on the API server. Keep it below the '--event-ttl' of the API server (one hour by default), or the event expires between repeats and the hold stops being visible, and keep 'spark.kubernetes.operator.events.minIntervalSeconds' within the rest of that TTL, that is at most the TTL minus this interval, or a repeat is dropped and the next one lands after the event expired. |
| spark.kubernetes.operator.reconciler.suspendHoldRequeueIntervalSeconds | Long | 1800 | true | Requeue interval (in seconds) at which a resource held by spec.suspend is reconciled, so that its SuspendHeld event is republished. A SparkApplication or SparkCluster suspended while running publishes no such event, since its Suspended status says so, and uses this interval only to look at its released resources again. The hold ends only when a user clears spec.suspend, which arrives as a watch event and reconciles right away, so nothing waits for this interval. It also paces a SparkApplication or SparkCluster which is suspended by a Kueue eviction and whose Workload is deactivated. That hold ends when the Workload is reactivated, which arrives as a watch event as well. It is deliberately much coarser than spark.kubernetes.operator.reconciler.intervalSeconds, since a suspended resource has nothing to observe while each republish costs a read and a write on the API server. Keep it below the '--event-ttl' of the API server (one hour by default), or the event expires between repeats and the hold stops being visible, and keep 'spark.kubernetes.operator.events.minIntervalSeconds' within the rest of that TTL, that is at most the TTL minus this interval, or a repeat is dropped and the next one lands after the event expired. |
| spark.kubernetes.operator.reconciler.terminationTimeoutSeconds | Integer | 30 | false | Grace period for operator shutdown before reconciliation threads are killed. |
| spark.kubernetes.operator.reconciler.trimStateTransitionHistoryEnabled | Boolean | true | true | When enabled, operator would trim state transition history when a new attempt starts, keeping previous attempt summary only. It also drops the history of a suspended SparkCluster when it is resumed. |
| spark.kubernetes.operator.watchedNamespaces | String | default | true | Comma-separated list of namespaces that the operator would be watching for Spark resources. If set to '*', operator would watch all namespaces. |
Expand Down
5 changes: 2 additions & 3 deletions docs/configuration.md
Original file line number Diff line number Diff line change
Expand Up @@ -102,7 +102,7 @@ replacing its `message`. A repeat carrying the same `message` may be paced by
The initial `Submitted` status of a new resource is not persisted to the API server on its own, so
no event is published for it. A `Submitted` event is published only when a `SparkApplication` or a
`SparkCluster` moves from `Suspended` back to `Submitted`, after `spec.suspend` is set back to
`false`, or after the eviction of a `SparkCluster` by Kueue.
`false`, or after its eviction by Kueue.

In addition, the operator publishes the following `Warning` events.

Expand All @@ -115,9 +115,8 @@ In addition, the operator publishes the following `Warning` events.
| `ClusterRequestFailed` | Applying the Services, StatefulSets, NetworkPolicy, HorizontalPodAutoscaler or PodDisruptionBudget of a `SparkCluster` fails in a way which may yet succeed, such as a throttled request or an internal server error from an admission webhook, so it is retried with the default interval instead of failing the cluster. Transport-level errors are skipped. |
| `SuspendReleaseFailed` | Deleting the driver pod (`SparkApplication`), the master and worker StatefulSets or the HorizontalPodAutoscaler or PodDisruptionBudget of the workers (`SparkCluster`), or the Kueue `Workload` of a resource suspended while running fails, or its pods cannot be listed. A failure which may clear on its own, such as a timeout, is skipped. Every failure is retried with the default interval. |
| `KueueResourceFlavorReadFailed` | Reading the `ResourceFlavor`s which Kueue assigned to an admitted `Workload` fails, so the node selector and the tolerations of the flavors cannot be applied. Transport-level errors are skipped and retried every 5 seconds, while a persistent failure is retried with the default interval. |
| `KueuePodsReadyUpdateFailed` | Recording the `PodsReady` condition on the Kueue `Workload` of a running `SparkApplication` or `SparkCluster` fails, e.g. without the permission for the `workloads/status` subresource, so Kueue may evict it by a `PodsReadyTimeout`, which releases a running `SparkCluster` and queues it again even if its master and workers are ready. Transport-level errors and conflicts are skipped. A `SparkCluster` retries them every 5 seconds and a persistent failure with the default interval, while a `SparkApplication` retries any failure with its next reconciliation, so that the recording does not hold back the observation of its driver. See [Kueue](spark_custom_resources.md#kueue). |
| `KueuePodsReadyUpdateFailed` | Recording the `PodsReady` condition on the Kueue `Workload` of a running `SparkApplication` or `SparkCluster` fails, e.g. without the permission for the `workloads/status` subresource, so Kueue may evict it by a `PodsReadyTimeout`, which releases it and queues it again, a `SparkApplication` with a new attempt, even if its pods are ready. Transport-level errors and conflicts are skipped. A `SparkCluster` retries them every 5 seconds and a persistent failure with the default interval, while a `SparkApplication` retries any failure with its next reconciliation, so that the recording does not hold back the observation of its driver. See [Kueue](spark_custom_resources.md#kueue). |
| `KueueDisabled` | A `SparkApplication` or `SparkCluster` labeled with `kueue.x-k8s.io/queue-name` is not queued, since `spark.kubernetes.operator.kueue.enabled` is disabled, so the driver (or master and worker) is requested without the Kueue admission. |
| `KueueEvictionIgnored` | Kueue evicted or deactivated the `Workload` of a `SparkApplication` whose driver is requested and whose attempt has neither stopped nor been suspended. The operator does not act on it, so the driver and executors keep running. An evicted `Workload` keeps holding its quota, while a deactivated one no longer counts against it. It is republished on every reconciliation while it lasts, subject to [`minIntervalSeconds`](#event-frequency). See [Kueue](spark_custom_resources.md#kueue). |

For a resource held by [`spec.suspend`](spark_custom_resources.md#suspend) or queued by
[Kueue](spark_custom_resources.md#kueue), the operator also publishes the following `Normal`
Expand Down
32 changes: 20 additions & 12 deletions docs/spark_custom_resources.md
Original file line number Diff line number Diff line change
Expand Up @@ -829,18 +829,26 @@ spec:
condition cannot be recorded, is admitted and evicted again and again until it is fixed,
suspended or deleted. With `blockAdmission`, it holds back every other workload from each
admission until its `Workload` is deleted after the backoff.
* Preemption is not honored yet after the driver of a `SparkApplication` is requested, and neither
is a `PodsReadyTimeout`, e.g. of an application which runs with fewer executors than
`spark.executor.instances`, or a deactivation. The operator checks the admission only before it
creates the resources, so after a later eviction the driver and executors keep running, and
until the attempt stops or is suspended, the `KueueEvictionIgnored`
[event](configuration.md#kubernetes-events) is published instead by default. An evicted
`Workload` keeps its Kueue quota until the driver and executors are released, so a workload
which preempts the application waits until then, while a deactivated one no longer counts
against it. Set `spec.suspend` to `true`, as described in [Suspend](#suspend), or delete the
application to release its quota earlier. With `blockAdmission`, an
application whose pods are not all ready holds back every other workload until they are ready
or its quota is released as described above.
* A `SparkApplication` whose driver is requested, i.e. from `DriverRequested` to
`RunningWithBelowThresholdExecutors`, is released like a running `SparkCluster` above when Kueue
evicts its `Workload`, e.g. to preempt it or by a `PodsReadyTimeout`, or deactivates it. The
application enters `Suspended` with a message giving the reason of the eviction, and the
operator deletes its driver pod, which deletes its executor pods. Like for a cluster, the
`Workload` is deleted once the pods are gone and its requeue backoff elapsed, or kept until it is
reactivated if it was deactivated, and then the application starts a new attempt from
`Submitted`, which is queued again with a new `Workload`. Like a resumed one, the new attempt
runs the application again from scratch and does not count against the restart limits, so it
starts even with `restartPolicy: Never`. An attempt whose driver has completed or failed by then
ends as usual instead, and `spec.suspend` takes precedence, so an application which is suspended
by it as well stays `Suspended` until it is resumed rather than being queued again.
* Unlike without Kueue, a `SparkApplication` whose driver and `spark.executor.instances` executors
are not all ready within the `waitForPodsReady` timeout, e.g. since some executors cannot be
scheduled, does not keep running with fewer executors, even if its `applicationTolerations`
allow it, e.g. in `RunningWithPartialCapacity`. Its `Workload` requested the quota for all of
them, so the `PodsReadyTimeout` releases the application as above, and it runs again from
scratch with a new attempt, which waits for the quota of all of them again. Like such a cluster,
an application whose executors never get ready is admitted and evicted again and again, with a
new attempt each time, until it is fixed, suspended or deleted.
* To limit the execution time of a `SparkApplication`, the Spark native `spark.driver.timeout` is
recommended instead of the `kueue.x-k8s.io/max-exec-time-seconds` label, which is not copied to
the `Workload`. It requires `spark.plugins=org.apache.spark.deploy.DriverTimeoutPlugin`, as in
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -246,6 +246,22 @@ public final class Constants {
"Application is resumed as spec.suspend is set to false, so a new attempt starts from "
+ "scratch.";

/**
* Message indicating that the application is evicted by Kueue after its driver was requested,
* which is followed by the reason and the message of the eviction. The operator tells such an
* application apart from one suspended by spec.suspend by this prefix of its persisted state
* message, so a changed wording no longer matches the applications suspended before.
*/
public static final String APP_EVICTED_MESSAGE =
"Application is evicted by Kueue, so its driver and executors are being released before its "
+ "Workload. It is queued again with a new attempt then, or once the Workload is "
+ "reactivated if it was deactivated.";

/** Message indicating that the application evicted by Kueue is submitted again. */
public static final String APP_REQUEUED_MESSAGE =
"Application is submitted again after its eviction by Kueue, so a new attempt starts from "
+ "scratch.";

// Spark Cluster Messages
/** Message indicating a failure to request the Spark cluster from the scheduler backend. */
public static final String CLUSTER_SCHEDULE_FAILURE_MESSAGE =
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -112,10 +112,10 @@ public enum ApplicationStateSummary implements BaseStateSummary {
TerminatedWithoutReleaseResources,

/**
* The application is suspended by spec.suspend after its driver was requested, so its driver
* and executors are released until it is resumed with a new attempt from Submitted. It is
* declared last, after the terminated states, so that the ordinal ranges of {@link
* #isStarting()} and {@link #isStopping()} exclude it.
* The application is suspended by spec.suspend, or by the eviction of its Kueue Workload, after
* its driver was requested, so its driver and executors are released until it is resumed with a
* new attempt from Submitted. It is declared last, after the terminated states, so that the
* ordinal ranges of {@link #isStarting()} and {@link #isStopping()} exclude it.
*/
Suspended;

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -268,13 +268,15 @@ public final class SparkOperatorConf {
* suspended while running publishes no such event, since its Suspended status says so, and uses
* this interval only to look at its released resources again. The hold ends only when a user
* clears {@code spec.suspend}, which arrives as a watch event and reconciles right away, so
* nothing waits for this interval. It is deliberately much coarser than {@link
* #RECONCILER_INTERVAL_SECONDS}, since a suspended resource has nothing to observe while each
* republish costs a read and a write on the API server. Keep it below the {@code --event-ttl} of
* the API server (one hour by default), or the event expires between repeats and the hold stops
* being visible, and keep {@link #KUBERNETES_EVENTS_MIN_INTERVAL_SECONDS} within the rest of
* that TTL, that is at most the TTL minus this interval, or a repeat is dropped and the next
* one lands after the event expired.
* nothing waits for this interval. It also paces a SparkApplication or SparkCluster which is
* suspended by a Kueue eviction and whose Workload is deactivated. That hold ends when the
* Workload is reactivated, which arrives as a watch event as well. It is deliberately much
* coarser than {@link #RECONCILER_INTERVAL_SECONDS}, since a suspended resource has nothing to
* observe while each republish costs a read and a write on the API server. Keep it below the
* {@code --event-ttl} of the API server (one hour by default), or the event expires between
* repeats and the hold stops being visible, and keep {@link
* #KUBERNETES_EVENTS_MIN_INTERVAL_SECONDS} within the rest of that TTL, that is at most the TTL
* minus this interval, or a repeat is dropped and the next one lands after the event expired.
*/
public static final ConfigOption<Long> SUSPEND_HOLD_REQUEUE_INTERVAL_SECONDS =
ConfigOption.<Long>builder()
Expand All @@ -288,7 +290,10 @@ public final class SparkOperatorConf {
+ "says so, and uses this interval only to look at its released resources "
+ "again. The hold ends only "
+ "when a user clears spec.suspend, which arrives as a watch event and "
+ "reconciles right away, so nothing waits for this interval. It is "
+ "reconciles right away, so nothing waits for this interval. It also paces a "
+ "SparkApplication or SparkCluster which is suspended by a Kueue eviction and "
+ "whose Workload is deactivated. That hold ends when the Workload is "
+ "reactivated, which arrives as a watch event as well. It is "
+ "deliberately much coarser than "
+ "spark.kubernetes.operator.reconciler.intervalSeconds, since a suspended "
+ "resource has nothing to observe while each republish costs a read and a "
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -923,7 +923,7 @@ public static boolean finishWorkload(
* timeout of Kueue does not evict it, and with `blockAdmission`, other workloads are admitted.
* Unlike the built-in Kueue integrations, the condition is not set back to `False` when a pod is
* lost later, e.g. an executor which Spark replaces, since an eviction would restart a running
* cluster with its applications, and is not acted on for a running application.
* cluster with its applications, or the attempt of a running application from scratch.
*
* <p>The informer cache answers whether the condition is still missing, so that a resource which
* reports it already, or whose pods are not ready yet, costs no request. Like {@link
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -55,7 +55,6 @@
import org.apache.spark.k8s.operator.reconciler.observers.AppDriverTimeoutObserver;
import org.apache.spark.k8s.operator.reconciler.reconcilesteps.AppCleanUpStep;
import org.apache.spark.k8s.operator.reconciler.reconcilesteps.AppInitStep;
import org.apache.spark.k8s.operator.reconciler.reconcilesteps.AppKueueEvictionStep;
import org.apache.spark.k8s.operator.reconciler.reconcilesteps.AppReconcileStep;
import org.apache.spark.k8s.operator.reconciler.reconcilesteps.AppResourceObserveStep;
import org.apache.spark.k8s.operator.reconciler.reconcilesteps.AppRunningStep;
Expand Down Expand Up @@ -207,7 +206,6 @@ protected List<AppReconcileStep> getReconcileSteps(final SparkApplication app) {
switch (app.getStatus().getCurrentState().getCurrentStateSummary()) {
case Submitted, ScheduledToRestart -> steps.add(new AppInitStep());
case DriverRequested, DriverStarted -> {
steps.add(new AppKueueEvictionStep());
steps.add(
new AppResourceObserveStep(
List.of(new AppDriverStartObserver(), new AppDriverReadyObserver())));
Expand All @@ -220,7 +218,6 @@ protected List<AppReconcileStep> getReconcileSteps(final SparkApplication app) {
RunningHealthy,
RunningWithPartialCapacity,
RunningWithBelowThresholdExecutors -> {
steps.add(new AppKueueEvictionStep());
steps.add(new AppRunningStep());
steps.add(new AppResourceObserveStep(List.of(new AppDriverRunningObserver())));
steps.add(new AppResourceObserveStep(List.of(new AppDriverTimeoutObserver())));
Expand Down
Loading