Repository navigation
Conversation
There was a problem hiding this comment.
Thanks for the PR, @TQJADE!
Both init steps now send a failed suspend check through the shared SuspendUtils.retryAfterCheckFailure. It publishes SuspendCheckFailed unless ReconcilerUtils.isTransientError holds, and keeps the retry at the default interval. That is what SPARK-59691 asks for. I checked a few things locally:
- With only the two call sites reverted, all 6 new test cases fail.
- The PR merges cleanly with current
main, andAppInitStepTestandClusterInitStepTeststill pass on the merge. getCurrentAttemptDriverPod()does throw on this path, since it does an APIGETonce the informer shows a live driver.
Minor
- 1. Cluster event message: The cluster passes
"master and workers", so the event reads "whether the master and workers of the suspended SparkCluster was requested". The check and the new doc row are about the master only. inline
| // look again with the steady-state interval instead. | ||
| log.warn("Failed to check whether the master of a suspended cluster exists.", e); | ||
| return completeAndDefaultRequeue(); | ||
| return SuspendUtils.retryAfterCheckFailure(context, e, "master and workers"); |
There was a problem hiding this comment.
Finding 1. With "master and workers" the message reads "Failed to check whether the master and workers of the suspended SparkCluster was requested". isMasterRequested reads only the master StatefulSet, and the new SuspendCheckFailed row in configuration.md says "the master of a suspended SparkCluster". Passing "master" fixes the grammar and matches both. The assertion at ClusterInitStepTest.java:899 would then check for "master of the suspended SparkCluster".
| return SuspendUtils.retryAfterCheckFailure(context, e, "master and workers"); | |
| return SuspendUtils.retryAfterCheckFailure(context, e, "master"); |
dongjoon-hyun
left a comment
There was a problem hiding this comment.
Thank you for the follow-up. I checked this PR on top of the current main. It merges cleanly, and spotlessCheck, checkstyleMain, checkstyleTest, pmdMain, AppInitStepTest and ClusterInitStepTest pass there. I left inline comments, which are summarized below in priority order.
- [P1]
ClusterInitStepL81: TheSparkClusterevent reads "the master and workers of the suspended SparkCluster was requested", which is ungrammatical and names workers that are never checked. Passing"master"fixes it. - [P2]
ClusterInitStepL83: A suspended cluster whose master exists looks it up again inholdForKueueAdmission, so a failure of that read is reported asKueueAdmissionRequestFailedwith a 5-second retry instead ofSuspendCheckFailedwith the default interval. - [P3]
SuspendUtilsL119: After the check recovers, the recoveringSuspendHeldcan be dropped byminIntervalSeconds, so the staleSuspendCheckFailedwarning stays the newest event for up to 30 minutes. - [P4]
ClusterInitStepL84: A cluster with a persisted history (lastKey() > 0) could go toSuspendedwithout the lookup, which leavesSuspendCheckFailedto the first attempt only. - [P5]
docs/configuration.mdL96 anddocs/spark_custom_resources.mdL547: The docs cover any failed check, but only aKubernetesClientExceptionis reported this way. - [P6]
SuspendUtilsL101-102:holdForSuspendstill swallows a failedWorkloadrelease without any warning. This is not introduced by this PR. - [P7]
AppInitStepTestL590-594: The new tests pin neither the full message nor its cause. - [P8]
AppInitStepTestL586-587: nit. The existingcaptureEvents(1)helper can be used.
P1, P2, P5, P7 and P8 are small changes within this PR. P3 and P4 are design suggestions, and P6 can be a separate JIRA issue.
| // look again with the steady-state interval instead. | ||
| log.warn("Failed to check whether the master of a suspended cluster exists.", e); | ||
| return completeAndDefaultRequeue(); | ||
| return SuspendUtils.retryAfterCheckFailure(context, e, "master and workers"); |
There was a problem hiding this comment.
[P1] With "master and workers", the event message renders as follows, and the log message starts the same way.
Failed to check whether the master and workers of the suspended SparkCluster was requested, will retry. KubernetesClientException: ...
The verb doesn't agree with the subject, and isMasterRequested reads only the master StatefulSet, so the workers are never checked. The new docs row, the Javadoc of REASON_SUSPEND_CHECK_FAILED and the sibling message at L249 (Failed to check whether the master exists ...) all say master, and the log message which this replaces (Failed to check whether the master of a suspended cluster exists.) was accurate as well.
| return SuspendUtils.retryAfterCheckFailure(context, e, "master and workers"); | |
| return SuspendUtils.retryAfterCheckFailure(context, e, "master"); |
ClusterInitStepTest L899 needs the same change ("master of the suspended SparkCluster").
| return completeAndDefaultRequeue(); | ||
| return SuspendUtils.retryAfterCheckFailure(context, e, "master and workers"); | ||
| } | ||
| if (!masterRequested) { |
There was a problem hiding this comment.
[P2] When this check finds the master, e.g. since the status update to RunningHealthy was lost, holdForKueueAdmission calls isMasterRequested again (L229 and L235). This double lookup predates this PR, but now that the lookup of a suspended cluster has its own reason, a failure of the second read is reported differently: as KueueAdmissionRequestFailed (Failed to check whether the master exists before requesting Kueue admission), with a 5-second retry of a transient failure, instead of SuspendCheckFailed with the default interval, as the new docs row says. In the rare case where the second read returns null, a suspended cluster would even request the Kueue admission. AppInitStep avoids this by keeping driverRequested outside the suspend branch and passing it to holdForKueueAdmission. Could we do the same here, i.e. masterRequested || isMasterRequested(context) at L229 and L235?
There was a problem hiding this comment.
added param isMasterRequested in the method holdForKueueAdmission
| + context.getResource().getKind() | ||
| + " was requested"; | ||
| log.warn("{}, will retry.", what, e); | ||
| if (!ReconcilerUtils.isTransientError(e)) { |
There was a problem hiding this comment.
[P3] Nothing retracts this warning once the check succeeds again. The recovering holdForSuspend publishes the same SuspendHeld message as before, which ConfigurableEventRecorder drops if it was published within minIntervalSeconds (300s), and the next one comes only after suspendHoldRequeueIntervalSeconds (1800s). For example, for a suspended SparkCluster on its first attempt, whose master is read from the API server on every reconciliation:
- t=0s: It is held, and
SuspendHeldis published. - t=60s: A spec edit triggers a reconciliation, the read fails with 403, and
SuspendCheckFailedis published. It is requeued after 120s. - t=180s: The read succeeds, and
SuspendHeldwith the same message is dropped. - t=1980s:
SuspendHeldis published again.
So, for about 30 minutes, the newest event says "will retry" although the check succeeded, while before this PR the newest event stayed the accurate SuspendHeld. The trigger is narrow, but since holdForSuspend already retracts a stale KueueAdmissionPending in its message, it may be worth handling here too, e.g. by saying so in the recovering SuspendHeld, or at least by documenting it.
| return SuspendUtils.retryAfterCheckFailure(context, e, "master and workers"); | ||
| } | ||
| if (!masterRequested) { | ||
| // Unlike a first attempt, a cluster resumed from Suspended has persisted this Submitted |
There was a problem hiding this comment.
[P4] For a cluster with a persisted history (lastKey() > 0, i.e. resumed from Suspended or requeued after an eviction), both outcomes of the lookup normally end in Suspended: a requested master completes the initialization and is suspended from RunningHealthy by ClusterSuspendStep, and an unrequested one goes to Suspended right here. Since ClusterSuspendStep deletes the StatefulSets by name and releases the Workload only after the pods are gone, whether or not the cluster went through RunningHealthy, we could check lastKey() > 0 before isMasterRequested and skip the lookup for such a cluster. Then the lookup and SuspendCheckFailed remain only for the first attempt, which is exactly the case without a persisted status that this PR targets. Currently, while the read keeps failing (e.g. 403), such a cluster republishes SuspendCheckFailed every 120s, and its persisted status keeps saying Cluster is resumed as spec.suspend is set to false., which is what this comment wants to avoid.
There was a problem hiding this comment.
| | `KueueAdmissionRequestFailed` | Creating, reading or deleting a stale Kueue `Workload` fails, or the driver or master cannot be read before the admission is requested. Transport-level errors are skipped and retried every 5 seconds, while a persistent failure is retried with the default interval. | | ||
| | `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 master and worker StatefulSets, the HorizontalPodAutoscaler or PodDisruptionBudget of the workers, or the Kueue `Workload` of a `SparkCluster` 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. | | ||
| | `SuspendCheckFailed` | Whether the driver of a suspended `SparkApplication`, or the master of a suspended `SparkCluster`, was requested cannot be checked, so the resource is neither held, which would release its Kueue `Workload`, nor started. Transport-level errors are skipped. Every failure is retried with the default interval. | |
There was a problem hiding this comment.
[P5] This row covers any failure to check, but only a KubernetesClientException reaches retryAfterCheckFailure. For a SparkApplication, a driver spec which cannot be built fails the same check differently. For example, when a suspended application has a live driver pod in the informer cache and its sparkConf was edited so that the spec no longer builds, SparkAppContext.getCurrentAttemptDriverPod (L117) throws a SparkException out of the reconciliation. It then gets ReconcileError on the first attempt only and the JOSDK retries, but neither SuspendCheckFailed nor the default interval, as suspendedAppWithCachedDriverAndUnbuildableSpecKeepsKueueWorkload asserts. How about narrowing it to a failed read, like the KueueAdmissionRequestFailed row?
| | `SuspendCheckFailed` | Whether the driver of a suspended `SparkApplication`, or the master of a suspended `SparkCluster`, was requested cannot be checked, so the resource is neither held, which would release its Kueue `Workload`, nor started. Transport-level errors are skipped. Every failure is retried with the default interval. | | |
| | `SuspendCheckFailed` | The driver of a suspended `SparkApplication`, or the master of a suspended `SparkCluster`, cannot be read to check whether it was requested, so the resource is neither held, which would release its Kueue `Workload`, nor started. Transport-level errors are skipped. Every such failure is retried with the default interval. | |
| server. An application held later, in `ScheduledToRestart`, keeps the status its previous attempt | ||
| already wrote and gets the same event, since that status says that a restart is due, not that the | ||
| next attempt is withheld. | ||
| next attempt is withheld. If the operator cannot check whether the driver (or master) was requested |
There was a problem hiding this comment.
[P5] Same as in configuration.md, a driver spec which cannot be built is not reported by SuspendCheckFailed.
| next attempt is withheld. If the operator cannot check whether the driver (or master) was requested | |
| next attempt is withheld. If the driver (or master) cannot be read to check whether it was requested |
| * server that is often the cause of it. Anything else is, since a suspended resource may have no | ||
| * persisted status to show it and would otherwise be retried without any signal. |
There was a problem hiding this comment.
[P6] This rationale applies to holdForSuspend above as well, which is not changed by this PR, so a follow-up JIRA issue is fine. KueueWorkloadUtils.releaseWorkload swallows a failed Workload deletion, and the hold then publishes a normal SuspendHeld without the release addendum and requeues after 30 minutes, without any warning. For example, if the Role of the operator lacks delete on workloads, or an admission webhook fails the deletion, the pending Workload of a suspended resource can be admitted into quota which nothing uses, while only the operator log says so. The dequeued path reports the same failure as KueueAdmissionRequestFailed (see failedReleaseOfPendingKueueWorkloadHoldsDequeuedCluster), and the reason given in the Javadoc of releaseWorkload, "the Workload is garbage collected with its owner anyway", doesn't hold for a suspended owner.
There was a problem hiding this comment.
| Assertions.assertTrue( | ||
| event.getValue().message().contains("driver of the suspended SparkApplication"), | ||
| event.getValue().message()); | ||
| Assertions.assertFalse( | ||
| event.getValue().message().contains("would not be requested"), event.getValue().message()); |
There was a problem hiding this comment.
[P7] These assertions pin neither the full message nor its cause. For instance, both new tests still pass when + EventUtils.describe(e) is dropped from retryAfterCheckFailure, and the contains check did not reveal the P1 message. Also, the assertFalse(...) can hardly fail, since the reason is asserted above and the template is fixed. How about comparing the full message, like the SuspendHeld tests in this file?
Assertions.assertEquals(
"Failed to check whether the driver of the suspended SparkApplication was requested, "
+ "will retry. KubernetesClientException: rejected",
event.getValue().message());The same applies to ClusterInitStepTest L898-902.
| ArgumentCaptor<EventRecord> event = ArgumentCaptor.forClass(EventRecord.class); | ||
| verify(eventRecorder).record(event.capture()); |
There was a problem hiding this comment.
[P8] nit. This file and ClusterInitStepTest already have captureEvents(int), which every other event assertion uses, including refusedDriverLookupBeforeKueueAdmissionPublishesEvent. EventRecord event = captureEvents(1).get(0); is equivalent and avoids the repeated event.getValue().
dongjoon-hyun
left a comment
There was a problem hiding this comment.
Thank you for the update. I checked the latest commit on top of the current main. It merges cleanly, and spotlessCheck, checkstyleMain, checkstyleTest, pmdMain, AppInitStepTest and ClusterInitStepTest pass there. I also tried a few mutations of the new masterRequested short-circuits. I left inline comments, which are summarized below in priority order.
- [P1]
ClusterInitStepL83-85: A suspended first attempt whose master exists still applies everything again and goes throughRunningHealthy, although the reasoning of the newlastKey() > 0branch holds for it as well. Moving it toSuspendedwould also make the newmasterRequestedparameter unnecessary. - [P2]
ClusterInitStepL237: No test covers thismasterRequested ||. All 555 tests of:spark-operator:testpass with it reverted. - [P3]
ClusterInitStepTestL455-457: Two assertions of the new queued test cannot fail for the regression which it targets, and thenullcase in its comment is not exercised. - [P4]
docs/spark_custom_resources.mdL551-552 anddocs/configuration.mdL96: ASuspendCheckFailedevent which is no longer refreshed has not cleared if the check now fails only at the transport level. - [P5]
ClusterInitStepL86-88: A master spec which cannot be built escapes the suspend check, so neither event is published and the pendingWorkloadis not released. This is not introduced by this PR. - [P6]
docs/spark_custom_resources.mdL549:SuspendCheckFailedis published only when events are enabled. - [P7]
ClusterInitStepL74-78: nit. A cluster queued again after an eviction does not say that it is resumed. - [P8]
ClusterInitStepL79-82: The PR description still says that the requeue behavior is unchanged. - [P9]
ClusterInitStepTestL411-413: The new queued test copies the stub block once more. Asuspendedparameter ofkueueFlavorsAreAppliedAgainWhenTheMasterExistscovers both L237 and L243. - [P10]
ClusterInitStepTestL328: nit. An unused stub contradicts thenever()verification below. - [P11]
ClusterInitStepTestL321: A requeued cluster, a transport-level failure andScheduledToRestartare not tested. - [P12]
docs/spark_custom_resources.mdL547-548: The Kueue section still says that suspending always releases theWorkload. - [P13]
docs/configuration.mdL96:up to 30 minutesholds only whileminIntervalSecondsis at most the suspend hold interval. - [P14]
SuspendUtilsL118-125: nit. This duplicatesKueueWorkloadUtils.retryAfterRequestFailure. - [P15]
AppInitStepTestL565-567: nit. The new parameterized tests duplicate the 503 tests next to them.
P1 is a design suggestion. If it is taken, P2, P3 and P9 are moot, since the masterRequested parameter goes away. P4, P6, P7, P8, P10, P12 and P13 are small changes within this PR, while P5, P11, P14 and P15 can be follow-ups.
| // 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. |
There was a problem hiding this comment.
[P1] The new comment above says that ClusterSuspendStep deletes the master and worker StatefulSets by name and releases the Workload only after the pods are gone, whether or not the master was requested. That holds for a first attempt whose master exists as well, yet such a cluster still completes its initialization here. It applies every resource again and persists RunningHealthy, and only then does ClusterSuspendStep move it to Suspended and delete what was just applied. For example, suppose the previous reconciliation applied the master StatefulSet but failed to apply the worker one with a retryable error, or the status update to RunningHealthy failed, and spec.suspend is set to true afterwards.
- The worker
StatefulSet, which never existed, is created, and its pods start only to be deleted in the next reconciliation. - If an admission webhook keeps failing the
NetworkPolicy, HPA or PDB with 500, onlyClusterRequestFailedrepeats. The suspension is never applied, while the master and its KueueWorkloadkeep the quota. - A non-retryable failure, e.g. 422 for a spec edited while suspended, or a node selector conflict of the Kueue flavors ends the cluster in
SchedulingFailure, although the user only suspended it.
Could we move such a cluster to Suspended as well, like the lastKey() > 0 branch? Then holdForKueueAdmission is reached only when spec.suspend is false, so the new masterRequested parameter and both masterRequested || are no longer needed. suspendAfterMasterRequestedCompletesInitialization and the new suspendedQueuedClusterWithRequestedMasterIsLookedUpOnce, which assert RunningHealthy for this case, would change accordingly. Since this path predates this PR, a follow-up JIRA issue is also fine.
| return KueueWorkloadUtils.handleDequeuedWorkload( | ||
| context, () -> masterRequested || isMasterRequested(context)); |
There was a problem hiding this comment.
[P2] No test covers this masterRequested ||. With this line reverted to () -> isMasterRequested(context), all 555 tests of :spark-operator:test still pass.
suspendedQueuedClusterWithRequestedMasterIsLookedUpOnceuses a cluster with the queue label, so it reaches only L243.kueueFlavorsAreAppliedAgainWhenTheMasterExists(queueLabelRemoved=true)is not suspended, somasterRequestedisfalsethere.suspendAfterMasterRequestedCompletesInitializationdoesn't stubgetCachedKueueWorkload(), sohandleDequeuedWorkloadreturns before it evaluates the supplier.
The uncovered case is a suspended first attempt whose queue label was removed after its admission. Without this line, a 403 of the second read is reported as KueueAdmissionRequestFailed (... before requesting Kueue admission), which is what the comment of the new test wants to avoid. A null applies the StatefulSets again without the flavors of the admitted Workload. P9 suggests a way to cover both L237 and L243.
| Assertions.assertEquals(ReconcileProgress.completeAndDefaultRequeue(), progress); | ||
| verify(mockClient.resource(masterStatefulSetSpec)).get(); | ||
| verify(mockClient, never()).resource(any(Workload.class)); |
There was a problem hiding this comment.
[P3] Two of these assertions cannot fail for the regression which this test targets. With L243 of ClusterInitStep reverted to if (isMasterRequested(context)), the second get() throws the stubbed 403, which holdForKueueAdmission turns into retryAfterRequestFailure. It returns completeAndDefaultRequeue() as well, and nothing reaches client.resource(workload). So only L456, L458 and L460 catch it, and the test passes with L455 and L457 alone.
The null case in the comment (a null would request the admission of a running master) is not exercised either. Even with .thenReturn(null) as the second answer, the test fails earlier, since the unstubbed priorityClasses().list() in KueueWorkloadPriority throws a NullPointerException, which ends in SchedulingFailure. How about saying in the comment that the get() count at L456 guards this, and either dropping L457 or covering the null case with priorityClasses().list() stubbed as well?
| the `SuspendHeld` event may follow only with its next repeat, so a `SuspendCheckFailed` event which | ||
| is no longer refreshed no longer applies. |
There was a problem hiding this comment.
[P4] This doesn't hold once the check fails only at the transport level, and neither does One which is no longer refreshed has cleared. in configuration.md. For example, the master read fails with 500, and SuspendCheckFailed is published. If the same outage then continues with 503, 504 or connection resets, retryAfterCheckFailure skips the event but still neither holds nor starts the cluster, so its pending Workload is not released either. A user who follows this sentence would conclude that the failure is over.
| the `SuspendHeld` event may follow only with its next repeat, so a `SuspendCheckFailed` event which | |
| is no longer refreshed no longer applies. | |
| the `SuspendHeld` event may follow only with its next repeat, so a `SuspendCheckFailed` event which | |
| is no longer refreshed no longer applies, unless the check now fails only at the transport level. |
| try { | ||
| masterRequested = isMasterRequested(context); | ||
| } catch (KubernetesClientException e) { |
There was a problem hiding this comment.
[P5] This catches only a KubernetesClientException, but isMasterRequested builds the whole master StatefulSet spec only for its name. If the spec cannot be built, e.g. after it was edited while suspended, the exception leaves the reconciliation, since the try at L106 is not reached yet. Then a suspended first attempt gets neither SuspendHeld nor SuspendCheckFailed, only a ReconcileError while the JOSDK retries it, and holdForSuspend never releases its pending Workload.
AppInitStep handles the same case with hasLiveDriverPod, and ClusterSuspendStep deletes the StatefulSets by name, since the spec may not even build after being edited while suspended. How about reading the master by name in isMasterRequested as well, i.e. apps().statefulSets().inNamespace(namespace).withName(getMasterStatefulSetName(name)).get()? This is not introduced by this PR, and I couldn't find a spec which passes the validation and still fails to build, so a follow-up JIRA issue is fine.
| @@ -321,17 +321,20 @@ void resumedClusterSuspendedAgainGoesBackToSuspended() { | |||
| new ClusterState(ClusterStateSummary.Submitted, Constants.CLUSTER_RESUMED_MESSAGE), | |||
There was a problem hiding this comment.
[P11] The new branches are tested only for this resumed history, a 503, and an application in Submitted.
- A cluster queued again after an eviction (
CLUSTER_REQUEUED_MESSAGE), which the new comment names, is not tested. The branch doesn't read the message today, butClusterSuspendSteptells the two apart, so a parameter over both messages would guard a later change. - A transport-level failure, e.g.
new KubernetesClientException("closed", new SocketException("reset")), is not tested for the suspend check of either step. OnlyReconcilerUtilsTestchecks its classification. SuspendCheckFailedof an application inScheduledToRestart, which the new docs sentence follows, is not tested either.
| next attempt is withheld. If the driver (or master) cannot be read to check whether it was requested | ||
| already, e.g. since the API server rejects the read, the resource is neither held nor started, and |
There was a problem hiding this comment.
[P12] The Kueue section still says that suspending a queued resource deletes its Workload to release the quota (L691), and the guide for disabling the integration asks to suspend the queued resources so that the operator releases their Workloads itself (L720-722). With this paragraph, that no longer holds while the check fails. For example, if the read keeps failing with 403, the Workload stays until the read succeeds, and a pending one may even be admitted meanwhile. The behavior predates this PR, but now that it is documented here, how about qualifying those two places too, e.g. by linking to this section?
| | `KueueAdmissionRequestFailed` | Creating, reading or deleting a stale Kueue `Workload` fails, or the driver or master cannot be read before the admission is requested. Transport-level errors are skipped and retried every 5 seconds, while a persistent failure is retried with the default interval. | | ||
| | `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 master and worker StatefulSets, the HorizontalPodAutoscaler or PodDisruptionBudget of the workers, or the Kueue `Workload` of a `SparkCluster` 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. | | ||
| | `SuspendCheckFailed` | The driver of a suspended `SparkApplication`, or the master of a suspended `SparkCluster`, cannot be read to check whether it was requested, so the resource is neither held, which would release its Kueue `Workload`, nor started. Transport-level errors are skipped. Every such failure is retried with the default interval. While the failure lasts, it is republished on every reconciliation, subject to [`minIntervalSeconds`](#event-frequency). Once the check succeeds again, the unchanged `SuspendHeld` event may be dropped by the same interval until its next repeat, so this event can stay the newest one for up to 30 minutes (`spark.kubernetes.operator.reconciler.suspendHoldRequeueIntervalSeconds`). One which is no longer refreshed has cleared. | |
There was a problem hiding this comment.
[P13] up to 30 minutes holds only while minIntervalSeconds is at most suspendHoldRequeueIntervalSeconds. With minIntervalSeconds=2400, which the Event frequency section uses as an example, the recovering SuspendHeld is dropped at its next repeat as well, so this event stays the newest one for about 60 minutes. The interval is also in the SuspendHeld row already. How about until its next repeat only? The suggestion also rewords the last sentence for P4.
| | `SuspendCheckFailed` | The driver of a suspended `SparkApplication`, or the master of a suspended `SparkCluster`, cannot be read to check whether it was requested, so the resource is neither held, which would release its Kueue `Workload`, nor started. Transport-level errors are skipped. Every such failure is retried with the default interval. While the failure lasts, it is republished on every reconciliation, subject to [`minIntervalSeconds`](#event-frequency). Once the check succeeds again, the unchanged `SuspendHeld` event may be dropped by the same interval until its next repeat, so this event can stay the newest one for up to 30 minutes (`spark.kubernetes.operator.reconciler.suspendHoldRequeueIntervalSeconds`). One which is no longer refreshed has cleared. | | |
| | `SuspendCheckFailed` | The driver of a suspended `SparkApplication`, or the master of a suspended `SparkCluster`, cannot be read to check whether it was requested, so the resource is neither held, which would release its Kueue `Workload`, nor started. Transport-level errors are skipped. Every such failure is retried with the default interval. While the failure lasts, it is republished on every reconciliation, subject to [`minIntervalSeconds`](#event-frequency). Once the check succeeds again, the unchanged `SuspendHeld` event may be dropped by the same interval until its next repeat, so this event can stay the newest one until then. One which is no longer refreshed has cleared, unless the check now fails only at the transport level. | |
| log.warn("{}, will retry.", what, e); | ||
| if (!ReconcilerUtils.isTransientError(e)) { | ||
| EventUtils.warn( | ||
| context.getEventRecorder(), | ||
| EventUtils.REASON_SUSPEND_CHECK_FAILED, | ||
| what + ", will retry. " + EventUtils.describe(e)); | ||
| } | ||
| return completeAndDefaultRequeue(); |
There was a problem hiding this comment.
[P14] nit. Except for the reason and the interval of a transient failure, this is the same as KueueWorkloadUtils.retryAfterRequestFailure: the same log line, the same what + ", will retry. " + EventUtils.describe(e) and the same gate. On top of the current main, the log.warn plus if (!ReconcilerUtils.isTransientError(e)) EventUtils.warn(...) pattern appears in seven places with this PR, e.g. ClusterSuspendStep, ClusterInitStep (ClusterRequestFailed), StatusRecorder and KueueWorkloadUtils.recordPodsReady. A helper like EventUtils.warnUnlessTransient(recorder, reason, message, e) would keep the policy in one place. This can be a follow-up.
| @ParameterizedTest | ||
| @ValueSource(ints = {403, 429, 500}) | ||
| void suspendedAppWithPersistentlyUnverifiableDriverPublishesEvent(int code) { |
There was a problem hiding this comment.
[P15] nit. This test and suspendedClusterWithPersistentlyUnverifiableMasterPublishesEvent differ from the 503 tests right above them only in the code and the event assertions. They could be one test each, like retryableFailureOfRequestingResourcesIsRetried(int code, int events) in ClusterInitStepTest, e.g. with @CsvSource({"503, 0", "403, 1", "429, 1", "500, 1"}). Then the Submitted assertion of the 503 tests would cover the other codes as well. The current split is also fine, since ClusterSuspendStepTest does the same.
|
Gentle ping, @TQJADE ~ |
dongjoon-hyun
left a comment
There was a problem hiding this comment.
Thank you for addressing round 2, @TQJADE. I checked the latest commit, 74076db, on top of the current main. It merges cleanly, and its CI passes except SparkOperatorConfigMapReconcilerTest.sanityTest in one JDK 26 job, which this PR does not touch. The new transition of a first attempt whose master exists looks right to me: ClusterSuspendStep releases such a cluster like any other suspended one. I left inline comments, which are summarized below in priority order.
- [P1]
ClusterInitStepL77: The PR description still saysThe requeue behavior is unchangedand doesn't list the newSuspendedtransitions, and its test section ends with an empty list. This is a follow-up of P8 of round 2. - [P2]
SuspendUtilsL48-49: The Javadoc ofholdForSuspendstill says that a resource requested before completes its initialization, which now holds forAppInitSteponly. - [P3]
ClusterInitStepL90:docs/architecture.mdL206 anddocs/spark_custom_resources.mdL588-589 don't show that a first attempt whose master exists now goes fromSubmittedtoSuspended. - [P4]
ClusterInitStepL85-88: A failed master lookup could go toSuspendedas well, since it is correct whether or not the master exists. If it stays as is, the comment could say why. - [P5]
ClusterInitStepL96-97: nit. This exit duplicates L78-79, and their comments frame the same outcome in opposite terms. - [P6]
ClusterInitStepL91-95: nit.ClusterSuspendStepkeeps the Services and the NetworkPolicy, soonly for ClusterSuspendStep to delete themoverstates it, and L220 leaves out a failed lookup.
P1, P2, P3, P5 and P6 are small changes within this PR. P4 is a design question, which can be a follow-up.
| // again. It goes back to Suspended instead, where ClusterSuspendStep deletes the master and | ||
| // worker StatefulSets by name and releases whatever it holds only after its pods are gone. | ||
| // That happens whether or not the master was requested, so it is not looked up. | ||
| if (cluster.getStatus().getStateTransitionHistory().lastKey() > 0) { |
There was a problem hiding this comment.
[P1] This is a follow-up of P8 of round 2. The new paragraph of the user-facing section covers part of it, but What changes were proposed in this pull request? still says the following, which becomes the commit message.
- The requeue behavior is unchanged: both cases are still retried with the default interval, and the resource is neither held nor started.
That no longer holds for a suspended SparkCluster which was resumed from Suspended or queued again after an eviction. Before this PR, a failed lookup kept it in Submitted and retried it with the default interval, while now it goes to Suspended here without the lookup. A first attempt whose master exists also goes to Suspended at L90-97 now, instead of completing its initialization to RunningHealthy. That section mentions neither change, and the new paragraph doesn't say that a resumed or requeued cluster is no longer looked up at all. In addition:
How was this patch tested?ends withPass the CIs with the newly added test cases:and nothing after it. It could list the two new parameterized tests, and say thatsuspendAfterMasterRequestedGoesToSuspendedreplacessuspendAfterMasterRequestedCompletesInitializationand thatresumedClusterSuspendedAgainGoesBackToSuspendednow asserts that the master is not looked up.- nit. The doc changes are not mentioned, and
spec.suspendand these events are not released yet in 1.0.0, which the template asks to clarify.
Could you update the description once more?
| @@ -88,4 +91,37 @@ static ReconcileProgress holdForSuspend(final BaseContext<?> context, final Stri | |||
| return completeAndRequeueAfter( | |||
There was a problem hiding this comment.
[P2] P1 of round 2 makes the Javadoc of this method stale. L48-49 still say:
* queued with outlives the Workload. Callers keep their own guard for resources requested
* before, which must complete their initialization instead of being held.
That now holds for AppInitStep only, since ClusterInitStep L90-97 sends a suspended first attempt whose master exists to Suspended rather than completing its initialization, as suspendAfterMasterRequestedGoesToSuspended asserts. The same commit already dropped the matching Like the suspend hold, a master requested before must complete its initialization from the Javadoc of ClusterInitStep.holdForKueueAdmission, so only this end is left. L48-49 are not in the diff, so I'm leaving this here. How about the following?
* queued with outlives the Workload. Callers keep their own guard for resources requested
* before, which are not held: an application whose driver was requested completes its
* initialization, while a cluster whose master was requested goes to Suspended instead.| context, statusRecorder, new ClusterState(Suspended, CLUSTER_SUSPENDED_MESSAGE)); | ||
| } | ||
| return SuspendUtils.holdForSuspend(context, "master and workers"); | ||
| if (masterRequested) { |
There was a problem hiding this comment.
[P3] This adds another path from Submitted to Suspended, which the docs don't show yet. Neither line below is in the diff, so I'm leaving this here.
docs/architecture.mdL206 labels the onlySubmitted --> Suspendededgespec.suspend=true after resume. A first attempt whose master already exists now takes it as well, from its unpersistedSubmittedstraight toSuspended, withoutRunningHealthy. The label also misses a cluster queued again after a Kueue eviction, which predates this PR.docs/spark_custom_resources.mdL588-589 say "Setting it totrueagain before the master is requested moves the cluster back toSuspended." Since L77-80, a resumed or requeued cluster goes back toSuspendedwhether or not the master was requested, sobefore the master is requestedis no longer the condition.
How about Submitted --> Suspended : spec.suspend=true after resume or Kueue requeue, or once the master is requested and "Setting it to true again moves the cluster back to Suspended."?
| // Whether the master is live is unknown, not answered. Holding would claim in an event | ||
| // that none was requested, and would release the Kueue quota of a running master, so | ||
| // look again with the steady-state interval instead. | ||
| log.warn("Failed to check whether the master of a suspended cluster exists.", e); | ||
| return completeAndDefaultRequeue(); | ||
| return SuspendUtils.retryAfterCheckFailure(context, e, "master"); |
There was a problem hiding this comment.
[P4] A design question, which can be a follow-up. After P1 of round 2, this lookup only decides between holdForSuspend and Suspended, and the comments at L72-76 and L91-95 explain that Suspended is correct whether or not the master was requested: ClusterSuspendStep deletes the StatefulSets by name, which is a no-op for a missing master, and releases the Workload only after the pods are gone. So a failed lookup could go to Suspended as well. That would persist a status, which kubectl get shows even though events are disabled by default, end the retry, and release a pending Workload instead of keeping it while it may be admitted (docs/spark_custom_resources.md L692-694).
On the other hand, a cluster which may never have started would then show Suspended (its master and workers are being released) instead of the unpersisted Submitted which the Suspend section describes, and isTransientError counts a single 429 or 500 as persistent, so one throttled read would leave it there until it is resumed. The same RBAC or API server failure also tends to fail the status patch and the deletions, and P4 of round 1 kept SuspendCheckFailed for exactly this first attempt. So I'm fine with keeping the current behavior. In that case, could the comment here say why such a cluster doesn't go to Suspended either, since the comments around it suggest that it could?
| return appendStateAndImmediateRequeue( | ||
| context, statusRecorder, new ClusterState(Suspended, CLUSTER_SUSPENDED_MESSAGE)); |
There was a problem hiding this comment.
[P5] nit. This exit is the same as L78-79, and their comments frame the same outcome in opposite terms: L72 starts with Unlike a first attempt, while L91-92 say that a first attempt whose master exists goes to Suspended as well. The rule is now that only a first attempt whose master is known to be absent is held. How about nesting the lookup under lastKey() == 0 and keeping a single transition with one merged comment? It behaves the same: a resumed or requeued cluster still skips the lookup, which resumedClusterSuspendedAgainGoesBackToSuspended checks, and only a KubernetesClientException is caught.
if (cluster.getSpec().isSuspend()) {
if (cluster.getStatus().getStateTransitionHistory().lastKey() == 0) {
final boolean masterRequested;
try {
masterRequested = isMasterRequested(context);
} catch (KubernetesClientException e) {
// (the current comment)
return SuspendUtils.retryAfterCheckFailure(context, e, "master");
}
if (!masterRequested) {
return SuspendUtils.holdForSuspend(context, "master and workers");
}
}
// (the comments of L72-76 and L91-95, merged)
return appendStateAndImmediateRequeue(
context, statusRecorder, new ClusterState(Suspended, CLUSTER_SUSPENDED_MESSAGE));
}| // A first attempt whose master StatefulSet already exists has been requested before, e.g. | ||
| // the status update to RunningHealthy failed. It goes to Suspended as well, rather than | ||
| // completing its initialization: that would apply every resource again only for | ||
| // ClusterSuspendStep to delete them, and a failure to apply them would keep the cluster | ||
| // from being suspended, or fail it, although it was only suspended. |
There was a problem hiding this comment.
[P6] nit. ClusterSuspendStep.releaseResources deletes only the StatefulSets and the HPA and PDB of the workers, and keeps the Services and the NetworkPolicy (The Services and the NetworkPolicy are kept, since they hold no quota and are applied again on resume.), so only for ClusterSuspendStep to delete them overstates it. The same goes for ClusterSuspendStep deletes what was applied in the comment of suspendAfterMasterRequestedGoesToSuspended. Also, since it is held or goes to Suspended first at L220 leaves out a failed lookup, which returns retryAfterCheckFailure at L88. How about the following, and e.g. since every suspend branch of reconcile returns first at L220?
| // A first attempt whose master StatefulSet already exists has been requested before, e.g. | |
| // the status update to RunningHealthy failed. It goes to Suspended as well, rather than | |
| // completing its initialization: that would apply every resource again only for | |
| // ClusterSuspendStep to delete them, and a failure to apply them would keep the cluster | |
| // from being suspended, or fail it, although it was only suspended. | |
| // A first attempt whose master StatefulSet already exists has been requested before, e.g. | |
| // the status update to RunningHealthy failed. It goes to Suspended as well, rather than | |
| // completing its initialization: that would apply every resource again only for | |
| // ClusterSuspendStep to release the master and workers, and a failure to apply them would | |
| // keep the cluster from being suspended, or fail it, although it was only suspended. |
dongjoon-hyun
left a comment
There was a problem hiding this comment.
Thank you for addressing round 3, @TQJADE. I checked the latest commit, 5864bac, on top of the current main. It merges cleanly, and its CI, spotlessCheck, checkstyleMain, checkstyleTest, pmdMain and the reconcilesteps tests pass there. I left inline comments, which are summarized below in priority order.
- [P1]
configuration.mdL96: TheexcludedReasonsrecipe, whichmainadded after the base of this PR, missesSuspendCheckFailed, although its message embedsEventUtils.describe(e)as well. - [P2]
configuration.mdL96:One which is no longer refreshed has cleareddoesn't hold for aSparkApplicationwhose driver spec stops building while the check fails. - [P3]
architecture.mdL206: The new label reads as if a requested master alone moved a cluster toSuspended. - [P4]
spark_custom_resources.mdL550: nit.A failure which may clear on its ownis ambiguous for 429 and 500, which publish the event. - [P5]
ClusterInitStepL81: nit.So look again with the steady-state interval instead.restates whatretryAfterCheckFailuredecides.
P14 and P15 of round 2 and P4 of round 3 are still open, which are fine as follow-ups.
| | `KueueAdmissionRequestFailed` | Creating, reading or deleting a stale Kueue `Workload` fails, or the driver or master cannot be read before the admission is requested. Transport-level errors are skipped and retried every 5 seconds, while a persistent failure is retried with the default interval. | | ||
| | `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 master and worker StatefulSets, the HorizontalPodAutoscaler or PodDisruptionBudget of the workers, or the Kueue `Workload` of a `SparkCluster` 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. | | ||
| | `SuspendCheckFailed` | The driver of a suspended `SparkApplication`, or the master of a suspended `SparkCluster`, cannot be read to check whether it was requested, so the resource is neither held, which would release its Kueue `Workload`, nor started. Transport-level errors are skipped. Every such failure is retried with the default interval. While the failure lasts, it is republished on every reconciliation, subject to [`minIntervalSeconds`](#event-frequency). Once the check succeeds again, the unchanged `SuspendHeld` event may be dropped by the same interval until its next repeat, so this event can stay the newest one until then. One which is no longer refreshed has cleared, unless the check now fails only at the transport level. | |
There was a problem hiding this comment.
[P1] After the base of this PR, main added the Event messages section (SPARK-59931), which recommends the following to keep the details of failures in the resource status and the operator log only.
spark.kubernetes.operator.events.excludedReasons=SchedulingFailure,ReconcileError,CleanupError,StatusUpdateFailed,ClusterRequestFailed,SuspendReleaseFailed,Kueue.*FailedSuspendCheckFailed embeds EventUtils.describe(e) as well, but matches none of them, since ConfigurableEventRecorder matches the whole reason (Pattern.matches(regex, reason)). So with the recommended setting, a 403 on the master read of a suspended SparkCluster still publishes the RBAC denial which names the ServiceAccount of the operator. With this PR, it is the only reason with EventUtils.describe(e) which the list misses. After rebasing, could you add SuspendCheckFailed to the list, or replace SuspendReleaseFailed with Suspend.*Failed, which doesn't match SuspendHeld?
| | `KueueAdmissionRequestFailed` | Creating, reading or deleting a stale Kueue `Workload` fails, or the driver or master cannot be read before the admission is requested. Transport-level errors are skipped and retried every 5 seconds, while a persistent failure is retried with the default interval. | | ||
| | `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 master and worker StatefulSets, the HorizontalPodAutoscaler or PodDisruptionBudget of the workers, or the Kueue `Workload` of a `SparkCluster` 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. | | ||
| | `SuspendCheckFailed` | The driver of a suspended `SparkApplication`, or the master of a suspended `SparkCluster`, cannot be read to check whether it was requested, so the resource is neither held, which would release its Kueue `Workload`, nor started. Transport-level errors are skipped. Every such failure is retried with the default interval. While the failure lasts, it is republished on every reconciliation, subject to [`minIntervalSeconds`](#event-frequency). Once the check succeeds again, the unchanged `SuspendHeld` event may be dropped by the same interval until its next repeat, so this event can stay the newest one until then. One which is no longer refreshed has cleared, unless the check now fails only at the transport level. | |
There was a problem hiding this comment.
[P2] This is a follow-up of P5 of round 1 and P4 of round 2. The row now covers a driver or master which cannot be read, but One which is no longer refreshed has cleared, unless the check now fails only at the transport level. doesn't hold for a SparkApplication whose check keeps failing for another reason. For example:
- A suspended application has a live driver pod in the informer cache, its read fails with 403, and
SuspendCheckFailedis published. - Its
sparkConfis then edited so that the driver spec no longer builds.SparkAppContext.getCurrentAttemptDriverPod()now throws aSparkExceptionat L117, which bypasses thecatch (KubernetesClientException e)ofAppInitStep. - The application gets only
ReconcileError, while it is still neither held nor started and keeps itsWorkload, if any.suspendedAppWithCachedDriverAndUnbuildableSpecKeepsKueueWorkloadcovers this path.
SuspendCheckFailed is no longer refreshed then, and with the excludedReasons which main recommends (see P1), ReconcileError is excluded, so the stale event reads as cleared. How about dropping the sentence, like spark_custom_resources.md, which says The operator log shows whether the check still fails.?
| RunningHealthy --> Suspended : spec.suspend=true or Kueue eviction | ||
| Suspended --> Submitted : spec.suspend=false or after Kueue eviction | ||
| Submitted --> Suspended : spec.suspend=true after resume | ||
| Submitted --> Suspended : spec.suspend=true after resume or Kueue requeue, or once the master is requested |
There was a problem hiding this comment.
[P3] This is a follow-up of P3 of round 3. or once the master is requested can be read as a trigger of its own, as if a cluster went from Submitted to Suspended when its master is requested, without spec.suspend. Besides, a cluster whose master is requested normally goes through RunningHealthy, since ClusterInitStep persists RunningHealthy in the reconciliation which applies the master. After this PR, a suspended cluster in Submitted goes to Suspended unless it is a first attempt without a master, so how about saying that?
| Submitted --> Suspended : spec.suspend=true after resume or Kueue requeue, or once the master is requested | |
| Submitted --> Suspended : spec.suspend=true, unless a first attempt has no master yet |
| next attempt is withheld. If the driver (or master) cannot be read to check whether it was requested | ||
| already, e.g. since the API server rejects the read, the resource is neither held nor started, and | ||
| the `SuspendCheckFailed` event is published instead when enabled, until the check succeeds. A | ||
| failure which may clear on its own, such as a timeout, is retried without the event. The event is |
There was a problem hiding this comment.
[P4] nit. A failure which may clear on its own follows the SuspendReleaseFailed sentence at L577, and the code means a transport-level failure by it (ReconcilerUtils.isTransientError). However, migration_guide.md in main uses the same words for a timeout, a 429, a 500 or a request which the API server never answered, while a throttled read (429) and an internal server error (500) do publish SuspendCheckFailed, as the new tests assert. How about naming the transport level, like the row in configuration.md?
| failure which may clear on its own, such as a timeout, is retried without the event. The event is | |
| failure at the transport level, such as a timeout, is retried without the event. The event is |
| // that none was requested, and would release the Kueue quota of a running master. | ||
| // Going to Suspended would show a cluster which may never have started as being | ||
| // released, and keep it there until it is resumed, even after a single throttled read. | ||
| // So look again with the steady-state interval instead. |
There was a problem hiding this comment.
[P5] nit. Now that SuspendUtils.retryAfterCheckFailure decides how to retry, this restates its @return (requeued after the steady-state interval), and would go stale if the helper changed, e.g. to retry a transport-level failure as soon as KueueWorkloadUtils.retryAfterRequestFailure does. The same goes for so look again with the steady-state interval instead. in AppInitStep L84-85.
| // So look again with the steady-state interval instead. |
|
Gentle ping, @TQJADE . |
dongjoon-hyun
left a comment
There was a problem hiding this comment.
Thank you for the follow-up. The overall direction looks good to me, and I didn't find a crash or a status corruption. I left a few comments, mainly about:
- Event under throttling: a
429still publishesSuspendCheckFailed, although the javadoc's own reason to skip transport-level failures (an event only adds load to an overloaded API server) applies to it as well. - First attempt with an existing master: the StatefulSet is looked up by name only, so a leftover of a deleted CR with the same name can move a never-started cluster to
Suspended. - Stale event after recovery: nothing supersedes
SuspendCheckFailedonce the check succeeds again. - Docs accuracy: after this PR, only a first attempt of a
SparkClusterlooks up its master, whileconfiguration.mdand the state diagram still read as if every suspended cluster were checked. - Minor: the
lastKey() == 0heuristic now gates more logic, the Kueue-requeued path has no test, and the warn-unless-transient pattern is duplicated once more.
| + context.getResource().getKind() | ||
| + " was requested"; | ||
| log.warn("{}, will retry.", what, e); | ||
| if (!ReconcilerUtils.isTransientError(e)) { |
There was a problem hiding this comment.
isTransientError returns false for 429, so a throttled read (API Priority and Fairness) publishes SuspendCheckFailed. The javadoc above skips transport-level failures because writing an event would only add load to an API server that is often the cause, which applies to 429 as well. The comment in ClusterInitStep even mentions a single throttled read. Shall we skip 429 here too?
| + " was requested"; | ||
| log.warn("{}, will retry.", what, e); | ||
| if (!ReconcilerUtils.isTransientError(e)) { | ||
| EventUtils.warn( |
There was a problem hiding this comment.
Once the check succeeds again, nothing supersedes this Warning. Since the following SuspendHeld can be dropped by minIntervalSeconds, a stale Failed to check ... will retry can stay the newest event in kubectl describe for up to the suspend hold interval (30 minutes by default), e.g. after an RBAC fix. The docs document it, but could we avoid it, e.g. by publishing a recovery or SuspendHeld event that is not deduplicated when the previous reconcile failed the check?
| + " of the suspended " | ||
| + context.getResource().getKind() | ||
| + " was requested"; | ||
| log.warn("{}, will retry.", what, e); |
There was a problem hiding this comment.
nit: The pattern of log.warn followed by EventUtils.warn unless ReconcilerUtils.isTransientError now exists in StatusRecorder, ClusterSuspendStep, ClusterInitStep, and here. A small helper such as EventUtils.warnUnlessTransient(recorder, reason, message, e) would keep the classification consistent, e.g. if 429 is reclassified. It can be a follow-up.
| // were requested, rather than applying every resource again. So the master of a resumed or | ||
| // requeued cluster is not looked up. | ||
| return appendStateAndImmediateRequeue( | ||
| context, statusRecorder, new ClusterState(Suspended, CLUSTER_SUSPENDED_MESSAGE)); |
There was a problem hiding this comment.
A first attempt whose master StatefulSet exists now goes to Suspended directly. Since isMasterRequested looks the StatefulSet up by name only, a StatefulSet left by a deleted SparkCluster with the same name, which is still being garbage-collected, makes a recreated suspend: true cluster Suspended (released), although it never ran, and it stays there until it is resumed. This is the outcome the comment in the catch block wants to avoid. Should we check the owner reference (UID) as well?
| if (cluster.getStatus().getStateTransitionHistory().lastKey() > 0) { | ||
| return appendStateAndImmediateRequeue( | ||
| context, statusRecorder, new ClusterState(Suspended, CLUSTER_SUSPENDED_MESSAGE)); | ||
| if (cluster.getStatus().getStateTransitionHistory().lastKey() == 0) { |
There was a problem hiding this comment.
The key of the state transition history, used as a proxy for a first attempt, now also decides whether the master is looked up at all. If a later change persists another Submitted state for a first attempt, e.g. while it waits for the Kueue admission, this branch flips silently. Could we name this condition, e.g. with a small helper or by checking the message of the Submitted state (CLUSTER_RESUMED_MESSAGE / CLUSTER_REQUEUED_MESSAGE), so that the intent is explicit?
| new ClusterState(ClusterStateSummary.Submitted, Constants.CLUSTER_RESUMED_MESSAGE), | ||
| true)); | ||
| cluster.getSpec().setSuspend(true); | ||
| // Not stubbed: a read would find a master, but the cluster must not look it up at all |
There was a problem hiding this comment.
This test covers a cluster resumed from Suspended (CLUSTER_RESUMED_MESSAGE) only. The PR description also mentions a cluster queued again after a Kueue eviction (CLUSTER_REQUEUED_MESSAGE). Could we parameterize this test with that message, ideally with a Kueue cluster, and verify that neither the master nor the Workload is touched?
| | `KueueAdmissionRequestFailed` | Creating, reading or deleting a stale Kueue `Workload` fails, or the driver or master cannot be read before the admission is requested. Transport-level errors are skipped and retried every 5 seconds, while a persistent failure is retried with the default interval. | | ||
| | `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 master and worker StatefulSets, the HorizontalPodAutoscaler or PodDisruptionBudget of the workers, or the Kueue `Workload` of a `SparkCluster` 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. | | ||
| | `SuspendCheckFailed` | The driver of a suspended `SparkApplication`, or the master of a suspended `SparkCluster`, cannot be read to check whether it was requested, so the resource is neither held, which would release its Kueue `Workload`, nor started. Transport-level errors are skipped. Every such failure is retried with the default interval. While the failure lasts, it is republished on every reconciliation, subject to [`minIntervalSeconds`](#event-frequency). Once the check succeeds again, the unchanged `SuspendHeld` event may be dropped by the same interval until its next repeat, so this event can stay the newest one until then. The operator log shows whether the check still fails. | |
There was a problem hiding this comment.
After this PR, only a first attempt of a SparkCluster looks up its master, while a resumed or requeued cluster goes to Suspended without the lookup and never gets this event. The same applies to the new paragraph in spark_custom_resources.md. How about saying the master of a suspended SparkCluster which has not started yet or similar?
| RunningHealthy --> Suspended : spec.suspend=true or Kueue eviction | ||
| Suspended --> Submitted : spec.suspend=false or after Kueue eviction | ||
| Submitted --> Suspended : spec.suspend=true after resume | ||
| Submitted --> Suspended : spec.suspend=true, unless a first attempt has no master yet |
There was a problem hiding this comment.
A first attempt whose master cannot be checked stays in Submitted as well, while this label reads as if it goes to Suspended. Maybe spec.suspend=true, unless a first attempt has no master yet or cannot check it?
# Conflicts: # docs/configuration.md # docs/spark_custom_resources.md # spark-operator/src/main/java/org/apache/spark/k8s/operator/reconciler/reconcilesteps/ClusterInitStep.java
There was a problem hiding this comment.
Thank you for the PR. I reviewed it and left 7 inline comments. A short summary:
Correctness
ConfigurableEventRecorder: aWarningresets the interval of theNormalevents of the resource, but not the other way round. OnceSuspendHeldsupersedesSuspendCheckFailed, a failure that comes back with the same message withinminIntervalSecondsis dropped. The newest event then says that the resource is held while the check is failing again.ClusterInitStep.isMasterRequested: the new owner check also changes the non-suspended Kueue path (holdForKueueAdmissionandhandleDequeuedWorkload). Neither the PR description nor the tests cover this.SuspendUtils.retryAfterCheckFailure: the PR description says that 429 is published as an event ("e.g. 403, 429 or 500"), but the code, the docs and the tests all skip it.
Design / Cleanup
isInitialSubmissiondepends on an implicit invariant (lastKey() == 0L), which a future change could break silently.- Every published warning scans the whole
lastRecordedmap withremoveIf. isOwnedByrepeats owner reference scanning which already exists elsewhere, e.g.KueueWorkloadUtils.controllerUid.- No test covers a failure with no response code.
Since suspend is not in the 1.0.0 release, I think no migration_guide.md item is needed.
| } | ||
| return previous; | ||
| }); | ||
| if (admitted.get() && event.type() == EventType.WARNING) { |
There was a problem hiding this comment.
This works in one direction only. A published Warning clears the Normal entries of the resource, but a published Normal does not clear its Warning entries. For example:
- t0: the lookup fails with 403 and
SuspendCheckFailedis published. - t1: the lookup succeeds and
SuspendHeldis published, since the warning cleared its entry. - t2 < t0 +
minIntervalSeconds: the lookup fails again with the same 403, so the same message is dropped as a repeat.
Then kubectl describe shows SuspendHeld as the newest event, while the check is failing and the Workload is kept. Should an admitted Normal also clear the Warning entries of the same resource, or should the check use the order of the last two events of the resource?
| // The name alone is not enough: a StatefulSet left by a deleted SparkCluster of the same | ||
| // name, which is still being garbage collected, was not requested by this one, and would | ||
| // otherwise move a recreated cluster which never ran to Suspended. | ||
| return master != null && isOwnedBy(master, context.getResource()); |
There was a problem hiding this comment.
isMasterRequested is also called by the non-suspended Kueue path: holdForKueueAdmission and handleDequeuedWorkload. So this owner check changes that path too. For example, a non-suspended queued SparkCluster reuses the name of a deleted one whose master StatefulSet is still being garbage collected. Before this PR, the leftover master counted as requested, so only applyAdmittedPodSetInfos ran. Now it counts as absent, so a new admission is requested, or handleDequeuedWorkload releases the pending Workload. This is probably the more correct behavior, but the PR description does not mention it and no test covers it without suspend. Could you add a test, and mention it in the PR description?
| + context.getResource().getKind() | ||
| + " was requested"; | ||
| log.warn("{}, will retry.", what, e); | ||
| if (!ReconcilerUtils.isTransientError(e) && e.getCode() != HTTP_TOO_MANY_REQUESTS) { |
There was a problem hiding this comment.
The PR description says that 429 is published ("Any other failure, e.g. 403, 429 or 500, is published as a SuspendCheckFailed Warning event"). But this line skips 429, and so do the docs and the tests (@ValueSource(ints = {503, 429}) expects no event). Could you update the PR description to match the code?
| * @return True if the cluster is in its initial submission, false otherwise. | ||
| */ | ||
| private static boolean isInitialSubmission(SparkCluster cluster) { | ||
| return cluster.getStatus().getStateTransitionHistory().lastKey() == 0L; |
There was a problem hiding this comment.
This check assumes that nothing persists a state before the master is requested. The Javadoc says so itself: "A change which persists another state before the master is requested must keep this check in step with it." Suppose a later change persists an intermediate state first, e.g. a Kueue pending status. Then a suspended first attempt would skip the master lookup and go to Suspended directly. ClusterSuspendStep would delete the StatefulSets by name, and a cluster which never started would be shown as released. Could we use a more explicit signal instead, e.g. whether the status was ever persisted?
| if (admitted.get() && event.type() == EventType.WARNING) { | ||
| lastRecorded | ||
| .keySet() | ||
| .removeIf(key -> key.type() == EventType.NORMAL && key.uid().equals(uid)); |
There was a problem hiding this comment.
Each admitted Warning scans the whole lastRecorded map, which holds the events of all resources in the last interval. During an API incident, many resources can publish warnings which are not transient, e.g. 403 or 500. Then the total cost grows quadratically, and it runs on the reconcile threads. Keying the state by uid, e.g. Map<String, Map<RecordedEvent, Long>>, or keeping a per-uid "last warning" timestamp which compute compares, would make this O(1).
| return master != null && isOwnedBy(master, context.getResource()); | ||
| } | ||
|
|
||
| private static boolean isOwnedBy(StatefulSet statefulSet, SparkCluster cluster) { |
There was a problem hiding this comment.
nit: This repeats the owner reference scan which already exists in KueueWorkloadUtils.controllerUid, and also in DriverResourceDecorator. How about a shared helper, e.g. in ModelUtils next to buildOwnerReferenceTo? Then a future fix, e.g. also requiring controller=true as controllerUid does, would be needed in one place only.
| @Test | ||
| void suspendedClusterWithUnverifiableMasterIsNotHeld() { | ||
| @ParameterizedTest | ||
| @ValueSource(ints = {503, 429}) |
There was a problem hiding this comment.
nit: The parameterized tests cover 503 and 429 (no event) and 403 and 500 (event). They do not cover a failure with no response code, e.g. a broken connection, which isTransientError classifies by its cause. A regression there would decide whether SuspendCheckFailed is written to an API server which is already failing. How about adding a NO_RESPONSE_CODE case with an IOException cause, both here and in AppInitStepTest?
What changes were proposed in this pull request?
This PR aims to publish a new
SuspendCheckFailedWarning event when the operator cannot check whether the driver of a suspendedSparkApplication, or the master of a suspendedSparkCluster, was requested already.EventUtils.REASON_SUSPEND_CHECK_FAILED.SuspendUtils.retryAfterCheckFailure, which bothAppInitStepandClusterInitStepnow call from their suspend branches when the lookup throws aKubernetesClientException. LikeClusterSuspendStepdoes forSuspendReleaseFailed, it classifies the failure withReconcilerUtils.isTransientError:SuspendCheckFailedWarning event.ClusterInitStepchanges two transitions of a suspendedSparkCluster:Suspended, or queued again after a Kueue eviction, no longer looks up its master and goes back toSuspendeddirectly. Previously, a failed lookup kept it inSubmittedand retried it with the default interval.ClusterSuspendStepdeletes the StatefulSets by name andreleases the
Workloadonly after the pods are gone, so the outcome does not depend on whether the master was requested.RunningHealthyfailed, now goes fromSubmittedtoSuspendedinstead of completing its initialization toRunningHealthy.Why are the changes needed?
This is a follow-up of SPARK-59680
Does this PR introduce any user-facing change?
Yes. When
spark.kubernetes.operator.events.enabledis set, a suspendedSparkApplicationorSparkClusterwhose driver or master lookup fails persistently now gets aSuspendCheckFailedWarning event. Previously no event was published. Transport-level failures behave as before.In addition, a suspended
SparkClusterwhich was resumed fromSuspended, queued again after an eviction, or whose master was already requested, now moves toSuspendedwithout completing its initialization toRunningHealthy.How was this patch tested?
Pass the CIs with the newly added test cases:
Was this patch authored or co-authored using generative AI tooling?
Generated-by: Claude Opus 5.5