[
https://issues.apache.org/jira/browse/KAFKA-20819?page=com.atlassian.jira.plugin.system.issuetabpanels:comment-tabpanel&focusedCommentId=18098377#comment-18098377
]
Guang Zhao commented on KAFKA-20819:
------------------------------------
Interesting. I think this issue also affects the trunk (not just 3.9.2) to a
degree -- though the trunk has bounded wait on `future.get` in
`TaskManager.waitForFuture`, it still makes the problem surface afterwards with
an obscure IllegalStateException on task not found. The remaining gap pointed
out above on detecting the silent death of the StateUpdater thread and error
reporting is still worth fixing on trunk, I think.
> 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)