[
https://issues.apache.org/jira/browse/KAFKA-20819?page=com.atlassian.jira.plugin.system.issuetabpanels:all-tabpanel
]
Matthias J. Sax resolved KAFKA-20819.
-------------------------------------
Resolution: Duplicate
> StreamThread hangs forever in TaskManager.waitForFuture() after its
> StateUpdater thread dies, silently dropping out of the consumer group
> -----------------------------------------------------------------------------------------------------------------------------------------
>
> Key: KAFKA-20819
> URL: https://issues.apache.org/jira/browse/KAFKA-20819
> Project: Kafka
> Issue Type: Bug
> Components: streams
> Affects Versions: 3.9.2
> Reporter: Stanislav Savulchik
> Priority: Major
>
> > *Note:* This report was drafted by an AI assistant (GitHub Copilot CLI)
> > based on
> > production evidence (broker logs, thread dump, per-thread application logs,
> > task
> > assignment snapshot) supplied by the reporter. All log excerpts and
> > line-number
> > references below were verified against the actual evidence files and the
> > Kafka 3.9
> > source tree; the reporter has reviewed and confirmed the analysis before
> > filing.
> h2. Summary
> A Kafka Streams {{StateUpdater}} thread can terminate on its own (its
> {{run()}} loop exits)
> without the owning {{StreamThread}}, {{TaskManager}}, or any supervisor ever
> finding out.
> {{DefaultStateUpdater.add()}}/{{remove()}} have no check for whether the
> thread is still
> alive: they unconditionally enqueue the request and, for {{remove()}}, return
> a
> {{CompletableFuture}} that only the {{StateUpdaterThread}} can complete. If
> that thread is
> already dead, the future is never completed.
> {{TaskManager.waitForFuture()}} ({{TaskManager.java}} lines ~700-720 on the
> 3.9 branch) then
> blocks on that future with a plain, *unbounded* {{future.get()}} (no timeout
> argument):
> {code:java}removedTaskResult = future.get();
> {code}
> This call happens on the {{KafkaConsumer.poll()}} thread, inside the
> rebalance listener
> callback ({{onPartitionsRevoked}} → {{TaskManager.handleRevocation()}} →
> {{revokeTasksInStateUpdater()}} → {{waitForFuture()}}). Because {{poll()}}
> never returns while
> blocked there, the consumer stops sending heartbeats/JoinGroup requests, and
> the broker's
> group coordinator eventually evicts the member on session-timeout — __while
> the JVM
> thread itself is still alive__, parked in {{WAITING}}, with no exception and
> no crash. Basic
> liveness checks (thread.isAlive(), process health checks) report the instance
> as healthy.
> h3. Environment / configuration
> * Kafka Streams 3.9.2
> * {{processing.guarantee}} = {{at_least_once}} (this is the intended/correct
> mode for this
> application; not exactly-once)
> * Topology deployed across 12 application instances, {{num.stream.threads}} =
> 2 per
> instance (24 StreamThreads total)
> * {{num.standby.replicas}} = 1
> * Default (KIP-892) {{DefaultStateUpdater}} enabled — one dedicated
> {{StateUpdater-N}} thread
> per {{StreamThread}}
> * Large RocksDB-backed state stores, standby/warmup replicas configured
> h3. Observed sequence (production evidence)
> # *2026-07-19 ~07:01–07:04* — a transient full network outage makes the broker
> cluster unreachable from one application instance ({{production-2}}). Both
> of its
> StreamThreads (and their restore consumers) see repeated {{Disconnecting
> from node …
> due to request/socket connection setup timeout}} across many broker nodes
> at the same
> time — this affects {{StreamThread-1}}/{{StateUpdater-1}} and
> {{StreamThread-2}}/
> {{StateUpdater-2}} simultaneously.
> # {{StateUpdater-1}} (belonging to the healthy {{StreamThread-1}}) keeps
> retrying broker
> connections for several more minutes and eventually recovers;
> {{StreamThread-1}} is
> healthy today.
> # {{StateUpdater-2}} (belonging to {{StreamThread-2}}) does *not* recover the
> same way.
> Its restore consumer logs {{Unsubscribed all topics or patterns and
> assigned
> partitions}} at {{07:04:51.883}}, and one second later, at
> {{07:04:52.933}},
> {{DefaultStateUpdater}} logs:
> {code} state-updater
> [sputnik-friend-locator-state-engine-production-2-StateUpdater-2] State
> updater thread stopped
> {code}
> This is the unconditional {{finally}}-block log line in
> {{DefaultStateUpdater$StateUpdaterThread.run()}}. Critically, *neither* of
> the two
> code paths that are supposed to precede it appear anywhere in the captured
> logs for
> this event:
> ** no {{"Shutting down state updater thread"}} (the log line the public
> {{shutdown(Duration)}} method emits before setting {{isRunning=false}}),
> and
> ** no {{"An unexpected error occurred within the state updater thread: …"}}
> (the
> {{ERROR}}-level line {{handleRuntimeException()}} emits before setting
> {{isRunning=false}}).
> This suggests the thread's {{run()}} loop exited via an uncaught
> {{Throwable}} that is
> _not_ a {{RuntimeException}} (e.g. an {{Error}}) — which bypasses the
> {{catch (RuntimeException)}} block entirely (so {{handleRuntimeException}}
> never runs and
> never logs), while the {{finally}} block still executes
> ({{clearInputQueue()}},
> {{clearUpdatingAndPausedTasks()}}, {{shutdownGate.countDown()}}, and the
> "State updater
> thread stopped" log). Such a {{Throwable}} propagating out of
> {{Thread.run()}} is handled
> by the default {{Thread.UncaughtExceptionHandler}}, which prints to
> {{System.err}} — a
> stream that was not captured in this evidence bundle (only stdout was
> captured), so we
> cannot pinpoint the exact {{Throwable}}. This gap is noted as a caveat;
> regardless of the
> trigger, the architectural problem below holds for any path that ends the
> {{StateUpdaterThread}} without going through {{shutdown(Duration)}}.
> # {{StateUpdater-2}} never restarts. No new thread by that name appears again.
> # Roughly *7 hours later*, at {{14:14:48}}, {{StreamThread-2}}
> (group-instance-id
> {{…production-2-2}}) performs a fresh {{JoinGroup}} and is briefly
> re-admitted to the
> consumer group (generation 28757, then a couple more generations through
> {{14:19:57}}).
> # At {{14:20:57}}, the broker {{GroupCoordinator}} removes that member again:
> {{"Member …production-2-2-e3dc728b… in group … has failed, removing it
> from the
> group"}} (heartbeat/session-timeout expiration). The group stabilizes one
> member
> short and *never re-admits it again* — it is still short one member days
> later.
> # A thread dump taken *~2 days later* (2026-07-21) still shows
> {{StreamThread-2}} parked
> at exactly this call chain:
> {code}
> java.util.concurrent.CompletableFuture.get(CompletableFuture.java:2072)
>
> org.apache.kafka.streams.processor.internals.TaskManager.waitForFuture(TaskManager.java:704)
>
> org.apache.kafka.streams.processor.internals.TaskManager.lambda$getNonFailedTasks$10(TaskManager.java:667)
> ...
>
> org.apache.kafka.streams.processor.internals.TaskManager.revokeTasksInStateUpdater(TaskManager.java:1215)
>
> org.apache.kafka.streams.processor.internals.TaskManager.handleRevocation(TaskManager.java:1114)
>
> org.apache.kafka.streams.processor.internals.StreamsRebalanceListener.onPartitionsRevoked(StreamsRebalanceListener.java:98)
>
> org.apache.kafka.clients.consumer.internals.ConsumerCoordinator.onJoinComplete(ConsumerCoordinator.java:402)
>
> org.apache.kafka.clients.consumer.internals.AbstractCoordinator.joinGroupIfNeeded(AbstractCoordinator.java:504)
>
> org.apache.kafka.clients.consumer.internals.AbstractCoordinator.ensureActiveGroup(AbstractCoordinator.java:415)
>
> org.apache.kafka.clients.consumer.internals.ConsumerCoordinator.poll(ConsumerCoordinator.java:511)
>
> org.apache.kafka.streams.processor.internals.StreamThread.pollRequests(StreamThread.java:1280)
> {code}
> i.e. it is stuck in the exact revoke path described above, waiting on a
> {{CompletableFuture}} from a {{StateUpdater-2}} thread that stopped
> existing 2 days
> earlier. There is no {{StateUpdater-2}} thread anywhere in that thread
> dump — only
> {{StateUpdater-1}}, which belongs to the healthy {{StreamThread-1}} and is
> legitimately
> {{RUNNABLE}}, actively restoring tasks.
> # Consequence: the client-side task-assignment view still lists
> {{StreamThread-2}} with
> only stale STANDBY/WARMUP tasks and *zero active tasks*, while the broker's
> consumer-group listing shows one fewer member than the number of live
> StreamThreads —
> a discrepancy that is very hard to diagnose operationally, since there is
> no error, no
> crash, and no metric distinguishing "healthy idle thread" from "livelocked
> thread that
> fell out of the group because its state updater died."
> h3. Relationship to existing tickets
> This shares the same underlying hazard — a {{StreamThread}} blocking on a
> {{StateUpdater}}-completed future during a rebalance listener callback — as:
> * KAFKA-19853 — *StreamThread blocks on StateUpdater during
> {{onAssignment()}}*, causing
> transaction timeouts under EOS. Same blocking pattern, different callback
> ({{onAssignment}} vs. {{onPartitionsRevoked}}) and different trigger (slow
> restore vs. a
> dead thread).
> * KAFKA-20456 (resolved) / KAFKA-20721 (open) — {{waitForFuture()}} timing
> out (after a
> *5-minute bounded* {{future.get(5, TimeUnit.MINUTES)}}, added for
> 4.1.2+/4.3.0) and
> producing an {{IllegalStateException}} ("ISE" —
> {{java.lang.IllegalStateException}}, thrown
> when a removal future resolves to {{null}}).
> * KAFKA-18355 / KAFKA-17789 — other cases of a Streams-internal thread
> blocking
> indefinitely / getting stuck.
> What appears to be new here, specific to 3.9.x:
> * {{TaskManager.waitForFuture()}} in 3.9.2 has **no timeout at all**
> (confirmed by reading
> the 3.9 branch source directly), so instead of eventually throwing the ISE
> described in
> KAFKA-20456/20721, the {{StreamThread}} hangs *forever*.
> * The trigger is not "the state updater is busy/slow" but that **the state
> updater
> thread is confirmed dead__ (absent from the thread dump, with an explicit
> "State
> updater thread stopped" log and no subsequent restart), and nothing in
> {{TaskManager}} or
> {{DefaultStateUpdater}} detects or reports that condition before a later
> caller blocks on
> a future that thread can no longer fulfill.
> h3. Suggested direction
> * Bound {{TaskManager.waitForFuture()}} (all call sites) with a timeout tied
> to
> {{max.poll.interval.ms}} / the rebalance timeout, and fail loudly (log +
> crash/replace the
> {{StreamThread}}) rather than allowing an indefinite park.
> * Have {{DefaultStateUpdater.add()}}/{{remove()}} check whether
> {{stateUpdaterThread}} is still
> alive (or fail the returned future immediately) instead of silently
> enqueuing work for a
> thread that has already exited.
> * Ensure {{StateUpdaterThread.run()}} cannot exit without going through a
> single,
> well-logged path: catch {{Throwable}} (not just {{RuntimeException}}) so
> unexpected {{Error}}s
> are logged at {{ERROR}} before the thread stops, and/or install a
> {{Thread.UncaughtExceptionHandler}} on the {{StateUpdaterThread}} so any
> uncaught fatal error
> is guaranteed to reach the structured application logs (not just
> {{System.err}}).
> * Consider a metric/log line when a {{StreamThread}} misses heartbeats due to
> a stuck
> rebalance-listener callback, to make this diagnosable without a thread dump.
> h2. Related issues
> KAFKA-19853, KAFKA-20456, KAFKA-20721, KAFKA-18355, KAFKA-17789
> h2. Attachments
> None at filing time. Evidence files (coordinator logs, task-assignment CSV,
> consumer-group state, thread dump, and the per-thread {{StateUpdater-2}} log)
> are available
> and can be attached later on request — recommend redacting
> instance/consumer-group names
> first, since this is a public ASF Jira.
--
This message was sent by Atlassian Jira
(v8.20.10#820010)