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;
This commit is contained in:
committed by
Artem Bilan
parent
5360aa3130
commit
0e0e8cebe3
@@ -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() {
|
||||
|
||||
@@ -397,7 +397,7 @@ public class KafkaMessageListenerContainer implements SmartLifecycle {
|
||||
}
|
||||
return;
|
||||
}
|
||||
} while (!hasErrors && !partitionsWithRemainingData.isEmpty());
|
||||
} while (!hasErrors && isRunning() && !partitionsWithRemainingData.isEmpty());
|
||||
}
|
||||
}
|
||||
if (wasInterrupted) {
|
||||
|
||||
@@ -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());
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user