diff --git a/spring-kafka/src/main/java/org/springframework/kafka/config/KafkaListenerEndpointRegistry.java b/spring-kafka/src/main/java/org/springframework/kafka/config/KafkaListenerEndpointRegistry.java index 3f37d0af..961afa49 100644 --- a/spring-kafka/src/main/java/org/springframework/kafka/config/KafkaListenerEndpointRegistry.java +++ b/spring-kafka/src/main/java/org/springframework/kafka/config/KafkaListenerEndpointRegistry.java @@ -1,5 +1,5 @@ /* - * Copyright 2014-2019 the original author or authors. + * Copyright 2014-2022 the original author or authors. * * Licensed under the Apache License, Version 2.0 (the "License"); * you may not use this file except in compliance with the License. @@ -227,14 +227,7 @@ public class KafkaListenerEndpointRegistry implements DisposableBean, SmartLifec @Override public void destroy() { for (MessageListenerContainer listenerContainer : getListenerContainers()) { - if (listenerContainer instanceof DisposableBean) { - try { - ((DisposableBean) listenerContainer).destroy(); - } - catch (Exception ex) { - this.logger.warn(ex, "Failed to destroy message listener container"); - } - } + listenerContainer.destroy(); } } diff --git a/spring-kafka/src/main/java/org/springframework/kafka/listener/MessageListenerContainer.java b/spring-kafka/src/main/java/org/springframework/kafka/listener/MessageListenerContainer.java index a756f97a..45627530 100644 --- a/spring-kafka/src/main/java/org/springframework/kafka/listener/MessageListenerContainer.java +++ b/spring-kafka/src/main/java/org/springframework/kafka/listener/MessageListenerContainer.java @@ -1,5 +1,5 @@ /* - * Copyright 2016-2019 the original author or authors. + * Copyright 2016-2022 the original author or authors. * * Licensed under the Apache License, Version 2.0 (the "License"); * you may not use this file except in compliance with the License. @@ -23,6 +23,7 @@ import org.apache.kafka.common.Metric; import org.apache.kafka.common.MetricName; import org.apache.kafka.common.TopicPartition; +import org.springframework.beans.factory.DisposableBean; import org.springframework.context.SmartLifecycle; import org.springframework.lang.Nullable; @@ -34,7 +35,7 @@ import org.springframework.lang.Nullable; * @author Gary Russell * @author Vladimir Tsanev */ -public interface MessageListenerContainer extends SmartLifecycle { +public interface MessageListenerContainer extends SmartLifecycle, DisposableBean { /** * Setup the message listener to use. Throws an {@link IllegalArgumentException} @@ -151,4 +152,9 @@ public interface MessageListenerContainer extends SmartLifecycle { throw new UnsupportedOperationException("This container does not support retrieving the listener id"); } + @Override + default void destroy() { + stop(); + } + } diff --git a/spring-kafka/src/test/java/org/springframework/kafka/listener/KafkaMessageListenerContainerTests.java b/spring-kafka/src/test/java/org/springframework/kafka/listener/KafkaMessageListenerContainerTests.java index 5fe8ff2b..78d8dfeb 100644 --- a/spring-kafka/src/test/java/org/springframework/kafka/listener/KafkaMessageListenerContainerTests.java +++ b/spring-kafka/src/test/java/org/springframework/kafka/listener/KafkaMessageListenerContainerTests.java @@ -1,5 +1,5 @@ /* - * Copyright 2016-2021 the original author or authors. + * Copyright 2016-2022 the original author or authors. * * Licensed under the Apache License, Version 2.0 (the "License"); * you may not use this file except in compliance with the License. @@ -643,8 +643,9 @@ public class KafkaMessageListenerContainerTests { inOrder.verify(consumer).commitSync(anyMap(), any()); inOrder.verify(messageListener).onMessage(any(ConsumerRecord.class)); inOrder.verify(consumer).commitSync(anyMap(), any()); - container.stop(); + container.destroy(); assertThat(advised).containsExactly("one", "two", "one", "two"); + assertThat(container.isRunning()).isFalse(); } @SuppressWarnings("unchecked")