Skip to content

[SPARK-59502] Support K8s Warning events for SparkApplication and SparkCluster failures - #826

Closed
TQJADE wants to merge 2 commits into
apache:mainfrom
TQJADE:event-record-oss
Closed

TQJADE wants to merge 2 commits into
apache:mainfrom
TQJADE:event-record-oss

Conversation

@TQJADE

@TQJADE TQJADE commented Sep 14, 2026 •

Copy link
Copy Markdown
Contributor

What changes were proposed in this pull request?

This PR aims to support publishing Kubernetes Warning events for SparkApplication and SparkCluster failures. The events are written into the namespace of the failed resource, so they show up under kubectl describe and kubectl get events.

The feature is disabled by default and is controlled by a new dynamic config, spark.kubernetes.operator.events.enabled.

Reason When
Failure state name A resource transitions into a failure state. For SparkApplication, these are SchedulingFailure, Failed, DriverEvicted, DriverStartTimedOut, DriverReadyTimedOut and ExecutorsStartTimedOut. For SparkCluster, these are SchedulingFailure and Failed.
ReconcileError A reconciliation throws. Only the first attempt of a failure episode publishes.
CleanupError A cleanup throws, so the resource cannot finish deleting.
StatusUpdateFailed A status patch is rejected. Transport-level errors (timeouts, 502, 503, 504) are skipped.

Implementation notes:

  • ConfigurableEventRecorder wraps the JOSDK default recorder and is registered once via ConfigurationServiceOverrider.withEventRecorder. It reads the config per event, so the dynamic override takes effect at runtime.
  • Failure-state events are emitted from StatusRecorder after a status patch succeeds and only when the status actually changed.
  • Each event uses its reason as the JOSDK event key, so repeated occurrences increment the count of one Event object.
  • BaseContext now holds the JOSDK context and provides getClient() and getEventRecorder().

Normal lifecycle events (for example, DriverRequested, RunningHealthy and Succeeded) are out of scope for this PR.

Why are the changes needed?

Most user-facing failures, such as a rejected driver pod or a driver start timeout, are recorded only in the custom resource status and the operator log. Publishing them as Warning events makes them visible through the standard Kubernetes tooling.

Does this PR introduce any user-facing change?

No behavior change by default. When spark.kubernetes.operator.events.enabled is true, the operator publishes the Warning events listed above. The existing Helm RBAC already grants the required permissions on events.

How was this patch tested?

Pass the CIs with the newly added unit tests. It was also manually tested on a Kubernetes cluster.

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

Generated-by: Claude Code

@dongjoon-hyun dongjoon-hyun left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

Thanks for the PR. I went through it against JOSDK 5.6.0 (DefaultEventRecorder / DefaultEventSink bytecode) and left inline comments. Two higher-level points that don't fit on a single line:

1. The chosen hook points rarely fire for the failures users actually see.
Every reconcile step converts failures into a status state instead of throwing: AppInitStep -> SchedulingFailure, AppValidateStep -> Failed, AppCleanUpStep -> ResourceReleased, AppDriverTimeoutObserver -> DriverStartTimedOut / DriverReadyTimedOut, and ClusterInitStep -> SchedulingFailure. So with events.enabled=true, an app rejected for quota/RBAC or timing out on driver start still shows no Events under kubectl describe; only operator-internal errors (NPEs, escaped KubernetesClientException from deleteResourceIfExists, status-patch failures) reach updateErrorStatus. The natural seam already exists: StatusRecorder fans out once per real state transition (after the patch succeeds, only when the status changed) and ApplicationStateSummary exposes isFailure() / isInfrastructureFailure(). Emitting there would cover the user-facing failures with built-in dedup. The persistStatus-failure hook is still useful as a complement.

2. Consider wiring the on/off flag once at the JOSDK level.
ConfigurationServiceOverrider.withEventRecorder(EventRecorder) exists in 5.6.0 and SparkOperator.overrideOperatorConfigs is already where every other framework toggle lives. A small delegating EventRecorder whose forContext() checks KUBERNETES_EVENTS_ENABLED.getValue() (keeping the dynamic toggle) would make Context.eventRecorder() itself safe and remove the Supplier indirection, the new BaseContext abstract method, and the static config read at each call site.

Minor: the PR title is truncated (... for SparkApplication and …). Since we squash-merge, it becomes the commit subject as-is.

retryInfo.isLastAttempt());
}
});
EventUtils.warn(

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.

JOSDK calls updateErrorStatus on every retry attempt (API_RETRY_MAX_ATTEMPTS defaults to 15, 5s x1.5 uncapped backoff, so a ~48+ min episode), and DefaultEventSink.emit is a synchronous GET + create/patch on the reconciler thread. This emits 15 times per failing resource, against an API server that is often the cause of the failure.

retryInfo.isLastAttempt() / getAttemptCount() is already in hand three lines above. Could we gate this inside the existing ifPresent block (e.g. first or last attempt only)? Same applies to SparkClusterReconciler.updateErrorStatus.

return;
}
try {
recorderSupplier.get().warn(reason, truncate(message));

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.

ResourceEventRecorder.warn(reason, message) builds an EventRecord with no key, and DefaultEventRecorder.eventName digests key().orElseGet(record::message) into the Event's metadata.name. So the full message is the dedup identity: any per-attempt variation (fabric8 KubernetesClientException messages embed the request URL and Status.toString() including retryAfterSeconds) creates a brand-new Event object per attempt instead of bumping count on one.

Suggest building the record explicitly with a stable key so repeats collapse into count increments:

recorderSupplier.get().record(
    EventRecord.builder()
        .type(EventType.WARNING)
        .reason(reason)
        .message(truncate(message))
        .key(reason)
        .build());

return true;
} catch (KubernetesClientException e) {
log.error("Error while persisting status to {}", newStatus, e);
EventUtils.warn(

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 catch is reached only after API_STATUS_PATCH_MAX_ATTEMPTS status patches already failed, i.e. the API server is rejecting or unreachable. The new call then issues two more blocking requests (GET + create/patch) to the same server on the reconciler thread, and AppReconcileStep.attemptStatusUpdate immediately requeues on false (bounded only by the per-resource rate limiter, ~5 loops / 15s). Cluster-side callers discard the boolean entirely.

Worth skipping the emit when the failure is transport-level (e.getCode() == 0 / ReconcilerUtils.isTransientError) or otherwise bounding it, so the observability path doesn't add load exactly when the control plane is degraded.

*/
public static void warn(
Supplier<ResourceEventRecorder> recorderSupplier, String reason, String message) {
if (!KUBERNETES_EVENTS_ENABLED.getValue()) {

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.

Two things sit outside the best-effort guard:

  1. Callers pass "... " + EventUtils.describe(e), so the cause walk and string build run before this flag check on every failure, even with the default enabled=false. The Supplier defers the cheap part (the recorder) but not the expensive part (the message).
  2. KUBERNETES_EVENTS_ENABLED is Boolean.class (not primitive), so ConfigOption.resolveValue goes through Jackson; readValue("null", Boolean.class) returns null, which is not caught, and !getValue() then NPEs. Because this line is outside the try, in the cleanup catch (RuntimeException e) { warn(...); throw e; } that NPE would replace the original cleanup failure, and in persistStatus it escapes the return false contract.

Low probability, but easy to close: take the message as Supplier<String> (or check the flag at the call sites) and move the flag read inside the try.

}

@Test
void warnSwallowsFailureFromRecorder() {

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.

Coverage gaps worth closing:

  • No test here verifies a publish on the success path; this one only asserts "no exception", so it passes even if the recorder.warn(...) call is removed. A verify(recorder).warn(eq(REASON_RECONCILE_ERROR), eq("boom")) would pin it.
  • truncate() has no test (and note it returns 1027 chars, substring(0, 1024) + "...", while the javadoc says 1024).
  • SparkClusterReconciler received the same four-site change with no new tests, and StatusRecorder.persistStatus's new site is never exercised (StatusRecorderTest uses mock(BaseContext.class)).
  • The new SparkAppReconcilerTest assertions use contains("Reconciliation failed.") / contains("cannot finish deleting"), which both the App and Cluster strings satisfy, so a kind copy-paste swap between the two reconcilers would go unnoticed. Asserting the full "Spark App ..." prefix would catch it.

*
* @return The event recorder.
*/
public abstract ResourceEventRecorder getEventRecorder();

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 abstract method plus the two identical overrides serve exactly one call site (StatusRecorder.persistStatus), while all four reconciler sites bypass it with context::eventRecorder on the JOSDK Context. That leaves two call conventions for the same value, glued together by the Supplier in EventUtils.

If we keep the util approach, two overloads warn(Context<?>, ...) / warn(BaseContext<?>, ...) that check the flag and then call eventRecorder() inline would drop the Supplier and this abstract method (or make it concrete over a shared josdkContext field instead of duplicating in both subclasses).

Also in EventUtils: DefaultEventRecorder.record already catches and logs emit failures, so the catch (RuntimeException) in warn only guards the supplier call; rootCauseOf's next.equals(rootCause) is unreachable because Throwable.getCause() returns null when cause == this, and the depth bound alone already guarantees termination; and truncate's null branch is unreachable since every caller concatenates a literal.


/** Utility class for publishing Kubernetes events about Spark resources. */
@Slf4j
public final class EventUtils {

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.

nit: this whole file is 4-space indented; the codebase (and CLAUDE.md, "Java with 2-space indentation") uses 2 spaces. Same for the new @AfterEach and three test methods in SparkAppReconcilerTest, and the EventUtils.warn( call sites use a 14-space hanging indent unlike the surrounding +4 continuations. Checkstyle has no Indentation module and Spotless runs no formatter, so CI won't flag it. Please reformat the added code (google-java-format style).

Related: the KUBERNETES_EVENTS_ENABLED description in SparkOperatorConf mixes trailing " + and leading + " fragments that split phrases ("into the" + " namespace "); the neighboring options use one + "..." fragment per line.

.appendNewStateAndPersist(any(SparkAppContext.class), any(ApplicationState.class));
}

@AfterEach

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.

nit: the PR flips the same flag through two different layers: SparkOperatorConfManager.INSTANCE.refresh(Map.of("spark.kubernetes.operator.events.enabled", ...)) here vs. TestUtils.setConfigKey (reflection on defaultValue) in EventUtilsTest. Since ConfigOption.resolveValue consults the override layer first, setConfigKey is shadowed whenever an override is live, and EventUtilsTest's @AfterEach can't restore it.

Suggest one idiom in both classes (the refresh(Map.of()) reset is the prevailing one in this repo), referencing the key via KUBERNETES_EVENTS_ENABLED.getKey() rather than a string literal. updateErrorStatusPublishesNoEventWhenDisabled also duplicates EventUtilsTest.warnDoesNothingWhenDisabled and could be dropped.

@dongjoon-hyun dongjoon-hyun changed the title [SPARK-59502] Support publishing Kubernetes Event objects for SparkApplication and … [SPARK-59502] Support publishing Kubernetes Event objects for SparkApplication and SparkCluster Sep 16, 2026
@dongjoon-hyun

Copy link
Copy Markdown
Member

Thank you for updating, @TQJADE . I fixed your PR title for you according to the JIRA information.

The PR code itself looks good to me now. Could you resolve the conflicts?

@TQJADE

TQJADE commented Sep 16, 2026

Copy link
Copy Markdown
Contributor Author

Thank you for updating, @TQJADE . I fixed your PR title for you according to the JIRA information.

The PR code itself looks good to me now. Could you resolve the conflicts?

Thanks a lot! I just resolved the conflicts. Kindly request your another around of review. @dongjoon-hyun

@dongjoon-hyun dongjoon-hyun changed the title [SPARK-59502] Support publishing Kubernetes Event objects for SparkApplication and SparkCluster [SPARK-59502] Support Kubernetes Warning events for SparkApplication and SparkCluster failures Sep 16, 2026
@dongjoon-hyun dongjoon-hyun changed the title [SPARK-59502] Support Kubernetes Warning events for SparkApplication and SparkCluster failures [SPARK-59502] Support K8s Warning events for SparkApplication and SparkCluster failures Sep 16, 2026
@dongjoon-hyun

dongjoon-hyun commented Sep 16, 2026 •

Copy link
Copy Markdown
Member

@TQJADE, I updated the PR title and description to match the actual scope of this PR.

The previous title, Support publishing Kubernetes Event objects ..., reads as general Kubernetes Events support. However, this PR publishes only Warning events for failures. Normal lifecycle events such as DriverRequested, RunningHealthy and Succeeded are not published. So I narrowed the title to:

[SPARK-59502] Support Kubernetes Warning events for SparkApplication and SparkCluster failures

I also rewrote the description based on the current code:

  • A table of the published reasons (failure state names, ReconcileError, CleanupError and StatusUpdateFailed) and when each one is emitted.
  • Brief implementation notes on ConfigurableEventRecorder, the failure-state events in StatusRecorder, and the reason-based event key.

@dongjoon-hyun dongjoon-hyun left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

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

+1, LGTM.

@dongjoon-hyun dongjoon-hyun added this to the 1.1.0 milestone Sep 16, 2026
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants