[SPARK-59907][SS] Release in-flight buffer when StreamingShuffleWriter.write() fails - #59177
Open
Jasraj-Dhanoa wants to merge 7 commits into
Open
Jasraj-Dhanoa wants to merge 7 commits into
Jasraj-Dhanoa wants to merge 7 commits into
Conversation
…-serialization failure StreamingShuffleWriter.write() takes a network buffer out of its shard and serializes a record into it before parking (putBuffer) or sending it. In that window the buffer is referenced only by a local variable, so if serialization throws, the task-completion cleanup (cleanupResources -> ShardState.cancel) cannot find it: the shard slot is already null and the buffer is not in the pool. It is therefore never released and leaks until GC reclaims the Unpooled direct buffer (a soft leak, worse under memory pressure). Wrap the serialize/checksum region in a try/catch that releases the buffer and returns its memory permit before rethrowing. The hand-off (putBuffer/send) is kept outside the guard because send() takes ownership and does its own cleanup. Adds a regression test that injects a serializer failing on the first value and asserts the captured buffer's refCnt returns to 0. NOTE: SPARK-XXXXX is a placeholder; replace with the real Jira id before merge. Co-authored-by: Isaac <no-reply@databricks.com>
…end pre-processing fails In ShardState.send(TimestampedBuffer) the serialization stream is closed, the checksum computed and the DataMessage built before the buffer is handed to the asynchronous send() whose completion callback owns the release/pool logic. If any of those synchronous steps throws, the buffer is referenced only by a local variable and leaks until GC, since cleanupResources() only reclaims pooled and shard-parked buffers (a soft leak, most likely under memory pressure while a buffer write grows off-heap memory). Guard that pre-hand-off region with a try/catch that releases the buffer and its memory permit before rethrowing. The send() call itself stays outside the guard, since that is where ownership legitimately transfers. Adds a regression test that injects a serializer whose stream close() throws and asserts the captured buffer's refCnt returns to 0. NOTE: SPARK-XXXXX is a placeholder; replace with the real Jira id before merge. Co-authored-by: Isaac <no-reply@databricks.com>
In ShardState.send(message, done), a failure while allocating or encoding the outgoing frame (e.g. an OutOfDirectMemoryError from the pooled allocator) released the frame and the message but never ran done(). For a DataMessage, done() is what releases the data buffer and returns its memory permit, so both leaked. The buffer had already been taken out of its shard, so the task-completion cleanup could not reclaim it either. Run done() on that path too, after releasing the message so the callback sees only the caller's reference, matching the other failure paths in this method. Adds a test that sends a DataMessage whose data buffer throws while being encoded and asserts done() runs exactly once. Co-authored-by: Isaac <no-reply@databricks.com>
Add a test that fails a task with an OutOfMemoryError mid-write, checks that its buffer is released, and runs a new shuffle on the same process. Also check that the mid-serialization failure test returns every semaphore permit, replace the SPARK-XXXXX placeholder with SPARK-59907, and trim the explanatory comments around the two error-path guards. Co-authored-by: Isaac <no-reply@databricks.com>
Keep only "Public for testing." to match the existing style in this file. Co-authored-by: Isaac <no-reply@databricks.com>
Jasraj-Dhanoa
marked this pull request as ready for review
October 1, 2026 01:07
Re-wrapping the row-size warning comment after indenting it into the try block dropped the words "more severe". Restore the original wording. Co-authored-by: Isaac <no-reply@databricks.com>
The recovery test called write() directly, which bypasses the executor's fatal-error handling. In a real executor an OutOfDirectMemoryError is fatal: Executor hands it to SparkUncaughtExceptionHandler, which exits the process, so that case cannot leave a lasting leak and the test's premise did not hold. Remove the test and the helpers only it used, drop the OutOfDirectMemoryError example from the encoding test's comment, and fix a comment that said a leaked buffer is reclaimed by GC (by default these buffers have no cleaner). Co-authored-by: Isaac <no-reply@databricks.com>
uros-b
reviewed
Oct 1, 2026
uros-b
left a comment
Member
There was a problem hiding this comment.
Thank you @Jasraj-Dhanoa for working on this. Seems like a surgical, idiom-correct leak fix: each error exit releases the locally-held ByteBuf and returns its permit, matching the writer's existing done()/cancel() reclamation, and the three injected-failure tests each map to one fixed path.
Could you please try to make the CI green? On first look, failed check seems like unrelated K8s-integration flake, but would be nice to clear if possible.
Also, let's ping @HeartSaVioR who has much more context here in streaming shuffle / RTM / SS
This branch has not been deployed
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?
StreamingShuffleWritercan leak an off-heap buffer when a write fails partway through. This PR releases the buffer on three error paths where it is referenced only by a local variable, or where the code that releases it is skipped:write(). The buffer has been taken out of its shard (or newly allocated) but not yet parked back or handed tosend(). If writing the key/value or computing the checksum throws, the buffer is now released and itsBUFFER_SIZEpermits are returned toallocatedBufferBytesSemaphorebefore the exception is rethrown.ShardState.send(timestampedBuffer). Closing the serialization stream, reading the data size, computing the checksum and building theDataMessageall run before the message is handed tosend(message, done), whose completion callback owns the release/pool logic. If any of these throw, the buffer is now released and its permits are returned before rethrowing.ShardState.send(message, done). Allocating the outgoing frame and encoding the message into it run before the frame is handed to the client. If either throws, the frame and the message were released but the completion callbackdone()was never run, so the data buffer and its permits leaked.done()is now run on this path too, after the message is released, matching the other failure paths in this method.The hand-off to
send(message, done)is intentionally left outside the first two guards: once it is called it owns the buffer, and releasing it again would corrupt the refcount.Why are the changes needed?
The task-completion listener (
cleanupResources()) only reclaims buffers that are pooled (bufferPool) or parked in a shard (ShardState.cancel()). A buffer that is referenced only by a local variable when an exception is thrown is invisible to it. By default, Netty allocates these buffers without a cleaner, so GC does not reclaim them either.These are defensive fixes. A leaked buffer only outlives the task when the failure is a non-fatal exception: a fatal error such as
OutOfDirectMemoryErrormakes the executor exit (Executorhands it toSparkUncaughtExceptionHandler), which frees its memory anyway. WithUnsafeRowSerializer, these paths are not expected to throw a non-fatal exception in normal operation; it would take a bug elsewhere (for example a refcount error). If one did occur, the buffer would stay allocated until the executor restarts, reducing the off-heap memory available to later tasks. Releasing the buffer on every error path makes buffer ownership consistent and keeps such a failure from turning into a silent leak.The file is byte-for-byte identical on
branch-4.3, so the change applies there as-is if a backport is wanted.Does this PR introduce any user-facing change?
No.
How was this patch tested?
Added three tests to
StreamingShuffleSuite:writer releases the in-flight buffer when serialization fails mid-write: a test serializer captures the buffer it writes into and throws on the first value. Asserts the captured buffer'srefCntdrops to 0 and all semaphore permits are returned.writer releases the in-flight buffer when the stream close fails on send: the same test serializer throws onclose()insideShardState.send(timestampedBuffer). Asserts the captured buffer'srefCntdrops to 0.shard send runs the completion callback when encoding fails: callsShardState.send(message, done)with aDataMessagewhose data buffer throws while being encoded. Assertsdone()runs exactly once and the message's reference is released.Each test fails without its fix (the buffer's
refCntstays at 1, one buffer's worth of permits is missing, ordone()is never run) and passes with it. All ofStreamingShuffleSuitepasses.Was this patch authored or co-authored using generative AI tooling?
Generated-by: Claude Code (Claude Opus)
This pull request and its description were written by Isaac.