Learn Labs
4. Kafka Consumers: Reading Data from Kafka

4.10 Exiting cleanly — `wakeup()` and `close()`

Skip close() and you pay session.timeout.ms of dead air on every deploy — across every partition that consumer owned.

consumer.wakeup() is THE ONLY consumer method safe to call from a DIFFERENT THREAD.

Mechanics:

"Calling wakeup will cause poll() to exit with WakeupException, or if consumer.wakeup() was called while the thread was not waiting on poll, the exception will be thrown on the next iteration when poll() is called. The WakeupException doesn't need to be handled, but before exiting the thread, you must call consumer.close()."

Why close() is not optional:

consumer.close():

  • COMMITS OFFSETS if needed
  • sends the group coordinator a message that the consumer is LEAVING → the coordinator triggers REBALANCING IMMEDIATELY

“You won't need to wait for the session to timeout before partitions from the consumer you are closing will be assigned to another consumer in the group.”

Skip close() and you pay session.timeout.ms of dead air on every deploy — across every partition that consumer owned.

Runtime.getRuntime().addShutdownHook(new Thread() {
    public void run() {
        System.out.println("Starting exit...");
        consumer.wakeup();                       // ① only safe cross-thread call
        try {
            mainThread.join();
        } catch (InterruptedException e) {
            e.printStackTrace();
        }
    }
});

...
Duration timeout = Duration.ofMillis(10000);     // ② a long poll timeout

try {
    // looping until ctrl-c, the shutdown hook will cleanup on exit
    while (true) {
        ConsumerRecords<String, String> records = movingAvg.consumer.poll(timeout);
        System.out.println(System.currentTimeMillis() + "--  waiting for data...");
        for (ConsumerRecord<String, String> record : records) {
            System.out.printf("offset = %d, key = %s, value = %s\n",
                record.offset(), record.key(), record.value());
        }
        for (TopicPartition tp: consumer.assignment())
            System.out.println("Committing offset at position:" +
                consumer.position(tp));
        movingAvg.consumer.commitSync();
    }
} catch (WakeupException e) {
    // ignore for shutdown                       // ③
} finally {
    consumer.close();                            // ④
    System.out.println("Closed consumer and we are done");
}

① "ShutdownHook runs in a separate thread, so the only safe action you can take is to call wakeup to break out of the poll loop."

② On long poll timeouts: "If the poll loop is short enough and you don't mind waiting a bit before exiting, you don't need to call wakeup — just checking an atomic boolean in each iteration would be enough. Long poll timeouts are useful when consuming low-throughput topics; this way, the client uses less CPU for constantly looping while the broker has no new data to return."

③ "You'll want to catch the exception to make sure your application doesn't exit unexpectedly, but there is no need to do anything with it."

④ "Before exiting the consumer, make sure you close it cleanly."