From 3fdf4682f2b72fbd8948323dcbc5079e21041eff 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) # Conflicts: # spring-kafka/src/main/java/org/springframework/kafka/listener/AbstractMessageListenerContainer.java --- .../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 f154c7b0..d23746a5 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 @@ -123,6 +123,8 @@ public abstract class AbstractMessageListenerContainer private volatile boolean running = false; + private volatile boolean fenced = false; + private volatile boolean paused; private volatile boolean stoppedNormally = true; @@ -274,6 +276,10 @@ public abstract class AbstractMessageListenerContainer return this.running; } + protected void setFenced(boolean fenced) { + this.fenced = fenced; + } + protected boolean isPaused() { return this.paused; } @@ -507,6 +513,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 483b0aef..0cf6543c 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 @@ -351,6 +351,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 bfff9e77..f59f318c 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; @@ -198,6 +199,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(); @@ -220,6 +222,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