[
https://issues.apache.org/jira/browse/KAFKA-20873?page=com.atlassian.jira.plugin.system.issuetabpanels:all-tabpanel
]
Matthias J. Sax updated KAFKA-20873:
------------------------------------
Fix Version/s: 4.4.1
(was: 4.4.0)
> KafkaStreams.removeStreamThread can deadlock with a concurrent REPLACE_THREAD
> uncaught-exception handler
> --------------------------------------------------------------------------------------------------------
>
> Key: KAFKA-20873
> URL: https://issues.apache.org/jira/browse/KAFKA-20873
> Project: Kafka
> Issue Type: Bug
> Components: streams
> Reporter: TengYao Chi
> Assignee: TengYao Chi
> Priority: Blocker
> Fix For: 4.2.2, 4.3.2, 4.4.1
>
>
> KafkaStreams.removeStreamThread(long) holds the changeThreadCount monitor
> across a StreamThread.waitOnThreadState(DEAD, ...) call. If the removal
> target happens to be the same StreamThread that concurrently entered the
> REPLACE_THREAD uncaught-exception handler, the whole client deadlocks and no
> further threads can be added, removed, or replaced.
> The affected KafkaStreams client stays alive (JVM/GC keeps ticking) but its
> Streams pipeline stops making progress indefinitely.
> The paths that collide
> KafkaStreams.removeStreamThread(long):
> synchronized (changeThreadCount) { //
> holds lock
> KafkaStreams.removeStreamThread(long) holds the changeThreadCount monitor
> across a StreamThread.waitOnThreadState(DEAD, ...) call. If the removal
> target happens to be the same StreamThread that concurrently entered the
> REPLACE_THREAD uncaught-exception handler, the whole client deadlocks and no
> further threads can be added, removed, or replaced.
> The affected KafkaStreams client stays alive (JVM/GC keeps ticking) but its
> Streams pipeline stops making progress indefinitely.
> The paths that collide
> ```
> KafkaStreams.removeStreamThread(long):
> synchronized (changeThreadCount) { //
> holds lock
> ...
> streamThread.shutdown(GroupMembershipOperation.LEAVE_GROUP); //
> signals target
> if (callingThreadIsNotCurrentStreamThread) {
> ...
> if (!streamThread.waitOnThreadState(State.DEAD, remainingTimeMs))
> { // BLOCKS while still holding the lock
> ...
> }
> }
> }
> KafkaStreams.addStreamThread() (also the tail of replaceStreamThread, which
> the REPLACE_THREAD uncaught-exception handler calls):
> synchronized (changeThreadCount) \{ ... }
> ```
> Deadlock scenario
> Two callers, one soak-style scheduleAtFixedRate loop that calls
> removeStreamThread(), and a StreamThread that throws an exception the
> uncaught-exception handler maps to REPLACE_THREAD:
> 1. The scheduler thread enters synchronized (changeThreadCount) inside
> removeStreamThread, iterates threads, picks StreamThread-N, calls
> streamThread.shutdown(...) (which only sets state to PENDING_SHUTDOWN), then
> blocks in waitOnThreadState(DEAD, ...) still holding the lock.
> 2. StreamThread-N's current runOnce invocation throws. The uncaught handler
> returns REPLACE_THREAD, which calls replaceStreamThread(...) ->
> deadThread.shutdown() -> addStreamThread(). That method also needs
> synchronized (changeThreadCount), so StreamThread-N blocks trying to
> reacquire the lock.
> 3. setState(DEAD) fires only at the end of completeShutdown(), which
> StreamThread-N can reach only by returning from the uncaught-exception
> handler and out of run(). It cannot, because it is blocked on
> changeThreadCount.
> 4. The scheduler thread waits forever for a DEAD transition that never comes.
> If the removal target is a different StreamThread than the one that hit the
> exception, no deadlock: that thread's completeShutdown does not touch
> changeThreadCount, so it reaches DEAD independently and notifies.
> Reproducer
> The Kafka Streams soak test in confluentinc/kafka-streams-soak-testing
> reproduces this reliably on kafka-clients 4.4.0-a1d9d8a517 (kip1071
> configuration, ALOS, self-hosted brokers): a periodic
> thread-add-remove-executor calls removeStreamThread() every 2h and an
> injected ThreadReplaceableException fires every 1-4h, so both paths race
> regularly. Multiple 15-30h hangs observed; representative thread dump:
> "thread-add-remove-executor" ... WAITING (on object monitor)
> at java.lang.Object.wait(Native Method)
> at
> o.a.k.streams.processor.internals.StreamThread.waitOnThreadState(StreamThread.java:1141)
> - locked <0x000000071352e4b8> (a java.lang.Object)
> at
> o.a.k.streams.KafkaStreams.removeStreamThread(KafkaStreams.java:1227)
> - locked <0x0000000712d8dee8> (a java.lang.Object)
> at
> o.a.k.streams.KafkaStreams.removeStreamThread(KafkaStreams.java:1211)
> "...-StreamThread-2" ... BLOCKED (on object monitor)
> at o.a.k.streams.KafkaStreams.addStreamThread(KafkaStreams.java:1140)
> - waiting to lock <0x0000000712d8dee8> (a java.lang.Object)
> at
> o.a.k.streams.KafkaStreams.replaceStreamThread(KafkaStreams.java:503)
> at
> o.a.k.streams.KafkaStreams.handleStreamsUncaughtException(KafkaStreams.java:522)
> at
> o.a.k.streams.KafkaStreams.lambda$setUncaughtExceptionHandler$0(KafkaStreams.java:475)
> at
> o.a.k.streams.processor.internals.StreamThread.run(StreamThread.java:959)
--
This message was sent by Atlassian Jira
(v8.20.10#820010)