Skip to content

KAFKA-20873: Fix deadlock between removeStreamThread and addStreamThread - #23004

Open
frankvicky wants to merge 1 commit into
apache:trunkfrom
frankvicky:KAFKA-20873-remove-stream-thread-deadlock
Open

KAFKA-20873: Fix deadlock between removeStreamThread and addStreamThread#23004
frankvicky wants to merge 1 commit into
apache:trunkfrom
frankvicky:KAFKA-20873-remove-stream-thread-deadlock

Conversation

@frankvicky

@frankvicky frankvicky commented Jul 31, 2026

Copy link
Copy Markdown
Contributor

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.

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.

`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>
@github-actions github-actions Bot added triage PRs from the community streams labels Jul 31, 2026
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

streams triage PRs from the community

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant