[SPARK-58973][SS] Fix state store batch read failure when StreamingQueryManager is not initialized - #58258
Open
shrirangmhalgi wants to merge 1 commit into
Conversation
…e StateStoreCoordinator RPC endpoint is not registered, instead of propagating SparkException.
shrirangmhalgi
commented
Aug 24, 2026
shrirangmhalgi
left a comment
Contributor
Author
There was a problem hiding this comment.
@HeartSaVioR / @anishshri-db / @Kimahriman could you please review. This fixes the batch state store read failure reported in #58211.
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
What changes were proposed in this pull request?
Make
StateStoreProvider.coordinatorRefgracefully returnNonewhen theStateStoreCoordinatorRPC endpoint is not available, instead of propagatingSparkException. This is done by wrapping theforExecutorcall in a try-catch that catchesSparkException(which wrapsRpcEndpointNotFoundExceptionfromawaitResult), logs a warning, and returnsNone.Why are the changes needed?
PR #50123 added
reportSnapshotUploadToCoordinator()toHDFSBackedStateStoreProvider.loadMap(), creating an unconditional dependency on theStateStoreCoordinatorfrom the state store read path. The coordinator endpoint is only registered whenStreamingQueryManageris instantiated — which is intentionally alazy valinSessionState(since SPARK-29423). In a fresh session that only does batch reads of state store data (e.g.,spark.read.format("statestore").load(path)), the endpoint doesn't exist and the read fails with[CANNOT_LOAD_STATE_STORE.UNCATEGORIZED].Does this PR introduce any user-facing change?
Yes.
spark.read.format("statestore").load(checkpointDir)now works in sessions that have never started a streaming query, without requiring the workaround of accessingspark.streamsfirst or settingspark.sql.streaming.stateStore.coordinatorReportSnapshotUploadLag=false.How was this patch tested?
Added
StateStoreCoordinatorBatchReadSuitewith two tests:StateStoreProvider.coordinatorRefreturnsNoneafter coordinator shutdownWas this patch authored or co-authored using generative AI tooling?
Yes. Co-Authored using Claude Opus 4.8
Closes #58211