From 0e0e8cebe386d46ffc876670d67b5fde654372bb Mon Sep 17 00:00:00 2001 From: Marius Bogoevici Date: Wed, 18 Mar 2015 11:41:06 -0400 Subject: [PATCH] INTEXT-153 Stop the KafkaMessageListenerContainer JIRA: https://jira.spring.io/browse/INTEXT-153 - ensure that the fetch loop is exited immediately after the listener container is stopped; - do not dispatch messages for processing once components are stopped; --- .../ConcurrentMessageListenerDispatcher.java | 4 ++- .../KafkaMessageListenerContainer.java | 2 +- .../QueueingMessageListenerInvoker.java | 32 ++++++++++--------- 3 files changed, 21 insertions(+), 17 deletions(-) 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()); + } } } }