KAFKA-20873: Fix deadlock between removeStreamThread and addStreamThread - #23004
Open
frankvicky wants to merge 1 commit into
Open
KAFKA-20873: Fix deadlock between removeStreamThread and addStreamThread#23004frankvicky wants to merge 1 commit into
frankvicky wants to merge 1 commit into
Conversation
`KafkaStreams.removeStreamThread(long)` used to hold the
`changeThreadCount` monitor across a `StreamThread.waitOnThreadState(DEAD, ...)`
call. When the removal target happened to be the same StreamThread that
concurrently entered the REPLACE_THREAD uncaught-exception handler, its
`replaceStreamThread -> addStreamThread` path blocked on the same lock,
so the target never reached `setState(DEAD)` (which only fires at the
end of `completeShutdown`), and the waiter waited forever.
Split `removeStreamThread` into three phases:
1. Pick a victim and signal shutdown under the lock.
2. Wait for the victim to reach DEAD *without* holding the lock.
3. Re-acquire the lock for bookkeeping (threads-list update, cache
resize, uncommitted-bytes resize).
Add `removeStreamThreadShouldReleaseLockBeforeWaitingForShutdown` — the
test parks the target's `waitOnThreadState` on a latch and verifies a
different thread can still acquire `changeThreadCount` while the removal
is waiting. Verified failing on unpatched code.
Co-Authored-By: Claude Opus 4.7 <noreply@anthropic.com>
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.
KafkaStreams.removeStreamThread(long)used to hold thechangeThreadCountmonitor across aStreamThread.waitOnThreadState(DEAD, ...)call. When the removaltarget happened to be the same StreamThread that concurrently entered
the REPLACE_THREAD uncaught-exception handler, its
replaceStreamThread -> addStreamThreadpath blocked on the same lock, so the target neverreached
setState(DEAD)(which only fires at the end ofcompleteShutdown), and the waiter waited forever.Split
removeStreamThreadinto three phases: 1. Pick a victim andsignal shutdown under the lock. 2. Wait for the victim to reach DEAD
without holding the lock. 3. Re-acquire the lock for bookkeeping
(threads-list update, cache resize, uncommitted-bytes resize).
Add
removeStreamThreadShouldReleaseLockBeforeWaitingForShutdown— thetest parks the target's
waitOnThreadStateon a latch and verifies adifferent thread can still acquire
changeThreadCountwhile the removalis waiting. Verified failing on unpatched code.
Delete this text and replace it with a detailed description of your
change. The PR title and body will become the squashed commit message.
If you would like to tag individuals, add some commentary, upload
images, or include other supplemental information that should not be
part of the eventual commit message, please use a separate comment.
If applicable, please include a summary of the testing strategy
(including rationale) for the proposed change. Unit and/or integration
tests are expected for any behavior change and system tests should be
considered for larger changes.