From 10f9726122accc02102c8a90e30e969ceb72d97a Mon Sep 17 00:00:00 2001 From: LokeshAlamuri <32193635+LokeshAlamuri@users.noreply.github.com> Date: Wed, 31 Jul 2024 22:28:38 +0530 Subject: [PATCH] GH-3371: Fence child containers after ConcurrentContainer stops Fixes: #3371 Containers are not restricted from starting after ConcurrentContainer stopped or restarted. These changes would fix this issue. * Enhancements to fence container after ConcurrentContainer stops (cherry picked from commit 20696f23900826b9c0ebe43042902659ffcb1a94) --- .../kafka/listener/AbstractMessageListenerContainer.java | 7 +++++++ .../kafka/listener/ConcurrentMessageListenerContainer.java | 1 + .../listener/ConcurrentMessageListenerContainerTests.java | 6 ++++++ 3 files changed, 14 insertions(+) diff --git a/spring-kafka/src/main/java/org/springframework/kafka/listener/AbstractMessageListenerContainer.java b/spring-kafka/src/main/java/org/springframework/kafka/listener/AbstractMessageListenerContainer.java index 8667f34c..8ef25815 100644 --- a/spring-kafka/src/main/java/org/springframework/kafka/listener/AbstractMessageListenerContainer.java +++ b/spring-kafka/src/main/java/org/springframework/kafka/listener/AbstractMessageListenerContainer.java @@ -124,6 +124,8 @@ public abstract class AbstractMessageListenerContainer private volatile boolean running = false; + private volatile boolean fenced = false; + private volatile boolean paused; private volatile boolean stoppedNormally = true; @@ -275,6 +277,10 @@ public abstract class AbstractMessageListenerContainer return this.running; } + protected void setFenced(boolean fenced) { + this.fenced = fenced; + } + @Deprecated(since = "3.2", forRemoval = true) protected boolean isPaused() { return this.paused; @@ -509,6 +515,7 @@ public abstract class AbstractMessageListenerContainer if (!isRunning()) { Assert.state(this.containerProperties.getMessageListener() instanceof GenericMessageListener, () -> "A " + GenericMessageListener.class.getName() + " implementation must be provided"); + Assert.state(!this.fenced, "Container Fenced. It is not allowed to start."); doStart(); } } diff --git a/spring-kafka/src/main/java/org/springframework/kafka/listener/ConcurrentMessageListenerContainer.java b/spring-kafka/src/main/java/org/springframework/kafka/listener/ConcurrentMessageListenerContainer.java index b8883ad9..4453696e 100644 --- a/spring-kafka/src/main/java/org/springframework/kafka/listener/ConcurrentMessageListenerContainer.java +++ b/spring-kafka/src/main/java/org/springframework/kafka/listener/ConcurrentMessageListenerContainer.java @@ -352,6 +352,7 @@ public class ConcurrentMessageListenerContainer extends AbstractMessageLis } } for (KafkaMessageListenerContainer container : this.containers) { + container.setFenced(true); if (container.isRunning()) { if (normal) { container.stop(() -> { diff --git a/spring-kafka/src/test/java/org/springframework/kafka/listener/ConcurrentMessageListenerContainerTests.java b/spring-kafka/src/test/java/org/springframework/kafka/listener/ConcurrentMessageListenerContainerTests.java index 3a8ba407..5561ae8a 100644 --- a/spring-kafka/src/test/java/org/springframework/kafka/listener/ConcurrentMessageListenerContainerTests.java +++ b/spring-kafka/src/test/java/org/springframework/kafka/listener/ConcurrentMessageListenerContainerTests.java @@ -17,6 +17,7 @@ package org.springframework.kafka.listener; import static org.assertj.core.api.Assertions.assertThat; +import static org.assertj.core.api.Assertions.assertThatExceptionOfType; import static org.mockito.ArgumentMatchers.any; import static org.mockito.ArgumentMatchers.anyString; import static org.mockito.BDDMockito.given; @@ -200,6 +201,7 @@ public class ConcurrentMessageListenerContainerTests { assertThat(container.metrics()).isNotNull(); Set> children = new HashSet<>(containers); assertThat(container.isInExpectedState()).isTrue(); + MessageListenerContainer childContainer = container.getContainers().get(0); container.getContainers().get(0).stopAbnormally(() -> { }); assertThat(container.isInExpectedState()).isFalse(); container.getContainers().get(0).start(); @@ -222,6 +224,10 @@ public class ConcurrentMessageListenerContainerTests { }); assertThat(overrides.get().getProperty(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG)).isNull(); this.logger.info("Stop auto"); + assertThat(childContainer.isRunning()).isFalse(); + assertThat(container.isRunning()).isFalse(); + // Fenced container. Throws exception + assertThatExceptionOfType(IllegalStateException.class).isThrownBy(() -> childContainer.start()); } @Test