[
https://issues.apache.org/jira/browse/KAFKA-20819?page=com.atlassian.jira.plugin.system.issuetabpanels:comment-tabpanel&focusedCommentId=18098453#comment-18098453
]
KAI-JAY TU commented on KAFKA-20819:
------------------------------------
Hi, I'd like to take a crack at this one.
I spent some time tracing through the code and the production evidence in the
ticket, and here's what I found:
*What's going wrong:*
StateUpdaterThread.run() only catches RuntimeException. When something like
OutOfMemoryError hits (as the reporter confirmed), it blows straight through to
the finally block. The thread logs "State updater thread stopped" and cleans up
queues, but nobody ever completes or fails the pending CompletableFutures
sitting in those queues.
On the other side, DefaultStateUpdater.remove() happily enqueues new work with
a fresh CompletableFuture — it never checks whether the thread is actually
alive. So when TaskManager.waitForFuture() calls future.get(5, MINUTES) on
trunk (or the unbounded future.get() on 3.9.x), it's waiting on a future that
will never complete.
Agreed with Guang's point — the trunk's 5-minute timeout mitigates the infinite
hang, but the underlying issue (silent StateUpdater death, undetected by
remove()/add()) still leads to the obscure IllegalStateException on task not
found. The root cause is the same.
*How I'd fix it (3 parts):*
# Widen the catch in StateUpdaterThread.run() to Throwable, so Errors get
properly logged and isRunning gets set to false. On thread exit, fail all
pending futures in the queue so nobody blocks on them.
# In remove() and add(), check isRunning() before enqueuing. If the thread is
already dead, complete the future exceptionally right away instead of queueing
work that'll never get processed.
# In TaskManager.waitForFuture(), don't just warn and return null on timeout —
check whether the state updater thread is actually dead, and if so, log at
ERROR with a clear message.
Before I go further — are you already working on a fix or planning to? Don't
want to duplicate effort. If not, I'm happy to pick this up. I've walked
through DefaultStateUpdater, TaskManager, and the rebalance listener chain, and
I can write the fix with tests that simulate a dead StateUpdaterThread. Happy
to put up a draft PR early for feedback.
One question: if the state updater dies, should we try to restart it, or is
failing fast (letting the StreamThread shut down) the right call? I'd lean
toward fail-fast to keep things simple, but want to hear your take.
> 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)