diff --git a/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/listener/ConcurrentMessageListenerDispatcher.java b/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/listener/ConcurrentMessageListenerDispatcher.java index c6b494db22..1f6c1b6987 100644 --- a/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/listener/ConcurrentMessageListenerDispatcher.java +++ b/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/listener/ConcurrentMessageListenerDispatcher.java @@ -113,7 +113,9 @@ class ConcurrentMessageListenerDispatcher implements Lifecycle { } public void dispatch(KafkaMessage message) { - delegates.get(message.getMetadata().getPartition()).enqueue(message); + if (isRunning()) { + delegates.get(message.getMetadata().getPartition()).enqueue(message); + } } private void initializeAndStartDispatching() { diff --git a/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/listener/KafkaMessageListenerContainer.java b/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/listener/KafkaMessageListenerContainer.java index 10be3063ee..3125af8064 100644 --- a/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/listener/KafkaMessageListenerContainer.java +++ b/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/listener/KafkaMessageListenerContainer.java @@ -397,7 +397,7 @@ public class KafkaMessageListenerContainer implements SmartLifecycle { } return; } - } while (!hasErrors && !partitionsWithRemainingData.isEmpty()); + } while (!hasErrors && isRunning() && !partitionsWithRemainingData.isEmpty()); } } if (wasInterrupted) { diff --git a/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/listener/QueueingMessageListenerInvoker.java b/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/listener/QueueingMessageListenerInvoker.java index 2f5c1d9fc4..56c3d2c67b 100644 --- a/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/listener/QueueingMessageListenerInvoker.java +++ b/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/listener/QueueingMessageListenerInvoker.java @@ -116,23 +116,25 @@ class QueueingMessageListenerInvoker implements Runnable, Lifecycle { while (this.running) { try { KafkaMessage message = messages.take(); - try { - if (messageListener != null) { - messageListener.onMessage(message); + if (isRunning()) { + try { + if (messageListener != null) { + messageListener.onMessage(message); + } + else { + acknowledgingMessageListener.onMessage(message, new DefaultAcknowledgment(offsetManager, message)); + } } - else { - acknowledgingMessageListener.onMessage(message, new DefaultAcknowledgment(offsetManager, message)); + catch (Exception e) { + if (errorHandler != null) { + errorHandler.handle(e, message); + } } - } - catch (Exception e) { - if (errorHandler != null) { - errorHandler.handle(e, message); - } - } - finally { - if (messageListener != null) { - offsetManager.updateOffset(message.getMetadata().getPartition(), - message.getMetadata().getNextOffset()); + finally { + if (messageListener != null) { + offsetManager.updateOffset(message.getMetadata().getPartition(), + message.getMetadata().getNextOffset()); + } } } }