diff --git a/spring-kafka/src/main/java/org/springframework/kafka/listener/KafkaMessageListenerContainer.java b/spring-kafka/src/main/java/org/springframework/kafka/listener/KafkaMessageListenerContainer.java index 9dd5f4c3..100c8fa0 100644 --- a/spring-kafka/src/main/java/org/springframework/kafka/listener/KafkaMessageListenerContainer.java +++ b/spring-kafka/src/main/java/org/springframework/kafka/listener/KafkaMessageListenerContainer.java @@ -53,6 +53,7 @@ import org.apache.kafka.clients.producer.Producer; import org.apache.kafka.common.Metric; import org.apache.kafka.common.MetricName; import org.apache.kafka.common.TopicPartition; +import org.apache.kafka.common.errors.AuthorizationException; import org.apache.kafka.common.errors.WakeupException; import org.springframework.context.ApplicationContext; @@ -908,6 +909,11 @@ public class KafkaMessageListenerContainer // NOSONAR line count ListenerConsumer.this.logger.error(nofpe, "No offset and no reset policy"); break; } + catch (AuthorizationException ae) { + this.fatalError = true; + ListenerConsumer.this.logger.error(ae, "Authorization Exception"); + break; + } catch (Exception e) { handleConsumerException(e); } @@ -1091,7 +1097,7 @@ public class KafkaMessageListenerContainer // NOSONAR line count } } else { - this.logger.error("No offset and no reset policy; stopping container"); + this.logger.error("Fatal consumer exception; stopping container"); KafkaMessageListenerContainer.this.stop(); } this.monitorTask.cancel(true); diff --git a/spring-kafka/src/test/java/org/springframework/kafka/listener/ConcurrentMessageListenerContainerMockTests.java b/spring-kafka/src/test/java/org/springframework/kafka/listener/ConcurrentMessageListenerContainerMockTests.java index 33d75389..e55e78eb 100644 --- a/spring-kafka/src/test/java/org/springframework/kafka/listener/ConcurrentMessageListenerContainerMockTests.java +++ b/spring-kafka/src/test/java/org/springframework/kafka/listener/ConcurrentMessageListenerContainerMockTests.java @@ -48,6 +48,7 @@ import org.apache.kafka.clients.consumer.ConsumerRecords; import org.apache.kafka.clients.consumer.OffsetAndTimestamp; import org.apache.kafka.clients.producer.Producer; import org.apache.kafka.common.TopicPartition; +import org.apache.kafka.common.errors.GroupAuthorizationException; import org.junit.jupiter.api.DisplayName; import org.junit.jupiter.api.Test; @@ -149,6 +150,35 @@ public class ConcurrentMessageListenerContainerMockTests { container.stop(); } + @SuppressWarnings({ "rawtypes", "unchecked" }) + @Test + void testConsumerExitWhenNotAuthorized() throws InterruptedException { + ConsumerFactory consumerFactory = mock(ConsumerFactory.class); + final Consumer consumer = mock(Consumer.class); + willAnswer(invocation -> { + Thread.sleep(100); + throw new GroupAuthorizationException("grp"); + }).given(consumer).poll(any()); + CountDownLatch latch = new CountDownLatch(1); + willAnswer(invocation -> { + latch.countDown(); + return null; + }).given(consumer).close(); + given(consumerFactory.createConsumer("grp", "", "-0", KafkaTestUtils.defaultPropertyOverrides())) + .willReturn(consumer); + ContainerProperties containerProperties = new ContainerProperties("foo"); + containerProperties.setGroupId("grp"); + containerProperties.setMessageListener((MessageListener) record -> { }); + containerProperties.setMissingTopicsFatal(false); + containerProperties.setShutdownTimeout(10); + ConcurrentMessageListenerContainer container = new ConcurrentMessageListenerContainer<>(consumerFactory, + containerProperties); + container.start(); + assertThat(latch.await(10, TimeUnit.SECONDS)).isTrue(); + verify(consumer).close(); + container.stop(); + } + @SuppressWarnings({ "rawtypes", "unchecked" }) @Test @DisplayName("Seek on TL callback when idle")