skip to content

How do you cleanly shut down a consumer that is blocked in poll()? Explain wakeup() vs close() and the typical shutdown-hook pattern.

level: seniorimportance: should knowfreq 60%

answer

  1. wakeup() = only thread-safe method
  2. poll() throws WakeupException
  3. catch it = clean stop, not error
  4. close() on owning thread, in finally
  5. close() = LeaveGroup -> fast rebalance

basics

~20 s

From another thread (e.g. a JVM shutdown hook) call consumer.wakeup() — the only thread-safe consumer method. It makes the in-progress poll() throw WakeupException. The poll thread catches it, breaks the loop, and calls close() itself (close() must run on the consumer's own thread).

solid answer

~40 s

Because KafkaConsumer is single-threaded, you can't stop it by calling close() from a control thread. Instead, the control thread (commonly a Runtime.getRuntime().addShutdownHook) calls consumer.wakeup(), which is the one thread-safe method. wakeup() causes the consumer's current or next poll() to immediately throw WakeupException. The poll-loop thread catches WakeupException, treats it as a normal stop signal (not an error), exits the loop, and then calls consumer.close() on its own thread. close() finalizes things: it sends a LeaveGroup so the coordinator rebalances promptly, optionally commits offsets, and releases connections. You typically join the consumer thread from the shutdown hook so the JVM doesn't exit before close() finishes. close() also accepts a timeout (close(Duration)) to bound how long final commits/leave can take.

go deeper

for a junior

Know that wakeup() stops a blocked poll() and you then close() the consumer.

for a middle

Explain the WakeupException catch-and-close pattern and why close() can't be called cross-thread.

for a senior

Detail the shutdown hook + join, close() sending LeaveGroup for fast rebalance, and bounding with close(Duration).

for a principal

Standardize graceful-shutdown across the consumer fleet: SIGTERM handling, drain semantics, final-commit guarantees, and bounded close timeouts under orchestrators like Kubernetes.

## The problem `KafkaConsumer` is not thread-safe and is confined to one thread, but shutdown signals (SIGTERM, app stop) arrive on *another* thread. You need a thread-safe way to interrupt a `poll()` that may be blocking for up to its timeout. ## wakeup() `consumer.wakeup()` is the **only** thread-safe `KafkaConsumer` method. Calling it sets an internal flag; the effect is that the consumer's **current `poll()` (or the next one, if none is active) throws `org.apache.kafka.common.errors.WakeupException`** immediately. It does not stop the consumer by itself — it just unblocks `poll()` so the owning thread can decide to stop. ## The canonical pattern ```java final Thread main = Thread.currentThread(); Runtime.getRuntime().addShutdownHook(new Thread(() -> { consumer.wakeup(); // safe from this other thread try { main.join(); } // wait for the loop to finish closing catch (InterruptedException ignored) {} })); try { while (true) { ConsumerRecords<String,String> records = consumer.poll(Duration.ofMillis(100)); process(records); consumer.commitSync(); } } catch (WakeupException e) { // expected on shutdown — swallow, do NOT rethrow } finally { consumer.close(); // runs on the owning thread } ``` Key points: catch `WakeupException` and treat it as a *clean stop*, not an error. Do the real cleanup (final commit if needed) and `close()` in the `finally`. The shutdown hook `join()`s the worker so the JVM waits for `close()` to complete. ## Why not just call close() from the hook? `close()` is **not** thread-safe. Invoking it from a thread other than the consumer's owner risks the `ConcurrentModificationException` guard or corrupt state. The correct division of labor: **wakeup() crosses threads; close() stays on the owning thread.** ## What close() does `close()` cleanly leaves the group by sending a **LeaveGroup** request, so the coordinator triggers a rebalance *immediately* instead of waiting for `session.timeout.ms` to expire — this minimizes the window where the consumer's partitions have no owner. It also: runs any auto-commit (if `enable.auto.commit=true`), invokes `ConsumerRebalanceListener.onPartitionsRevoked` for a final commit opportunity, and releases sockets/threads. `close(Duration timeout)` bounds how long these final operations may take before forcing shutdown; `close()` with no arg uses a default (30s). Always close in a `finally` so connections aren't leaked. ## Subtleties - A `WakeupException` can also surface from other blocking calls (e.g. `commitSync`) if `wakeup()` was called — handle it around the loop, not just the `poll()`. - Distinguish from `Thread.interrupt()`: interrupting the consumer thread raises `InterruptException`, a different signal; `wakeup()` is the intended, supported mechanism. - If you call `wakeup()` and there's no active poll, the *next* poll throws — the flag is sticky until consumed.

  • Why call close() instead of just letting the process exit?
    close() sends a LeaveGroup so the coordinator rebalances the consumer's partitions immediately rather than after session.timeout.ms; it also performs a final auto-commit / onPartitionsRevoked commit and releases connections. Skipping it leaves partitions unowned for up to the session timeout and may leak resources.
  • Why must the shutdown hook join the consumer thread?
    wakeup() returns immediately; the loop still needs time to catch WakeupException, commit, and close(). Joining the worker keeps the JVM alive until that finishes, so you don't half-close mid-commit and lose the clean LeaveGroup.
  • What if WakeupException is thrown from commitSync() rather than poll()?
    It can — any blocking consumer call observes a pending wakeup. So scope your try/catch around the whole loop body, treat WakeupException as the stop signal wherever it appears, and still close() in finally.

saying these in an interview costs you the question

  • Calling close() from the shutdown-hook thread directly
  • Treating WakeupException as an error to log/retry
  • Using Thread.interrupt() instead of wakeup()
  • Forgetting to close(), leaving partitions unowned until session.timeout.ms

context