I am trying to handle shutdown events in Kafka Consumers. API docs : https://kafka.apache.org/21/javadoc/index.html?org/apache/kafka/common/errors/WakeupException.html
public class KafkaConsumerRunner implements Runnable {
private final AtomicBoolean closed = new AtomicBoolean(false);
private final KafkaConsumer consumer;
public void run() {
try {
consumer.subscribe(Arrays.asList("topic"));
while (!closed.get()) {
ConsumerRecords records = consumer.poll(Duration.ofMillis(10000));
// Handle new records
}
} catch (WakeupException e) {
// Ignore exception if closing
if (!closed.get()) throw e;
} finally {
consumer.close();
}
}
// Shutdown hook which can be called from a separate thread
public void shutdown() {
closed.set(true);
consumer.wakeup();
}
}
Here why does the document expect us the handle if the boolean is false? if (!closed.get()) throw e; If an external thread calls shutdown closed will be set to true always. Is there any use-case I am missing?