Repository navigation
[fix][client] Complete a v5 send in place while its client is closing - #26686
Merged
merlimat merged 1 commit intoSep 22, 2026
Merged
Conversation
A send that fails because the client is closing has its failure delivered on the client's completion executor. The client shuts that executor down with shutdownNow(), which drops whatever is still queued on it, so a completion handed over in the window between the client entering the Closing state and the shutdown could be lost: the caller's future then never completes. finish() already completes in place once the executors are gone (RejectedExecutionException); do the same while the client is closing, before they are gone. The new test holds the client's only completion executor, closes the client, issues a send once the client is closing and checks that the send's future still fails with AlreadyClosedException. Without the fix the future never completes.
merlimat
approved these changes
Sep 22, 2026
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.
Fixes #26685
Main Issue: #26685
Motivation
A v5 send that fails because its client is closing can leave the caller's future pending forever.
ScalableTopicProducer.finish()hands the failure to the client's completion executor, andPulsarClientImpl.shutdown()stops that executor withshutdownNow(), which drops the tasks still queued on it. A completion queued in the window between the client enteringClosingand the executor shutdown is lost. The existingRejectedExecutionExceptionfallback only covers completions submitted after the executor is gone.This is what made
V5ProducerSegmentGoneTest#closingTheClientWhileSendsAreInFlightFailsThemWithoutRecreatingProducerstime out on a loaded CI runner (twice in a row, including the TestNG retry): the test keeps issuing sends while the client closes, and on that runner some completions were still queued when the executor was shut down. See #26685. Found while running the CI for #26672 on my fork, whose change does not touch the v5 client; #26684 came out of the same CI runs.Modifications
ScalableTopicProducer.finish(): complete the future in place when the client is closing (client.v4Client().isClosed(), the same checkisShuttingDown()already uses), the same way it already does once the executor rejects the task. No other behaviour changes; results that arrive while the client is open are still delivered on the completion executor.V5ProducerClosingCompletionTest: holds the client's only completion executor (callbackThreads(1)), closes the client on another thread, issues a send once the client reports closing and asserts that the send fails withAlreadyClosedExceptioninstead of never completing. Without the fix the future never completes (3/3 runs:TimeoutExceptionafter 5 s); with it the test passes (3/3).Verifying this change
This change added tests and can be verified as follows:
V5ProducerClosingCompletionTest#aSendFailedWhileTheClientIsClosingCompletesEvenIfItsCompletionWasQueuedfails on master and passes with this change.V5ProducerSegmentGoneTestandV5ProducerBackpressureTeststill pass (12 tests across the three classes), plus checkstyle and spotless forpulsar-client-v5and the broker tests.Does this pull request potentially affect one of the following parts:
If the box was checked, please highlight the changes
While the client is closing, a failed send's future is completed on the thread that failed it instead of on the completion executor. That is the thread
finish()already uses once the executor is gone.A question for reviewers
The check uses the underlying client's state,
client.v4Client().isClosed(), because that is what this class already does inisShuttingDown()(andDagWatchClientdoes the same). The alternative I considered is forPulsarClientV5to record that closing has begun itself, a flag set inclose(),closeAsync()andshutdown()and exposed to the package asisClosed(), so that the v5 layer stops reading the v4 client's state;isShuttingDown()would then use it too. I have that variant ready and can switch if you would rather the v5 client own this. I kept the smaller change here so the bug fix does not carry a design decision with it.Documentation
doc-not-neededMatching PR in forked repository
PR in forked repository: namest504#6