From 91d09c8ad3fcc481286e2ac37753d2aede0a4139 Mon Sep 17 00:00:00 2001 From: Artem Bilan Date: Thu, 29 Jul 2021 16:39:57 -0400 Subject: [PATCH] Add deprecation suppression for spring-kafka-2.8 Spring for Apache Kafka 2.8 has introduced an new `CommonErrorHandler` and deprecated its retying components including `RetryingMessageListenerAdapter`. * For proper compatibility with `spring-kafka-2.8.0`, which is going to be a foundation for upcoming Spring Boot 2.6, it is better to suppress deprecations and don't raise such a concern to end-users. In the future we will revise retrying logic to expected behavior from `spring-kafka-2.8.0`. Or will do home-made one as it is now with `AmqpInboundChannelAdapter`, for example. See https://github.com/spring-projects/spring-integration/issues/3605 --- .../kafka/inbound/KafkaInboundGateway.java | 7 ++++--- .../KafkaMessageDrivenChannelAdapter.java | 16 ++++++++-------- 2 files changed, 12 insertions(+), 11 deletions(-) diff --git a/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/inbound/KafkaInboundGateway.java b/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/inbound/KafkaInboundGateway.java index 3d502cbe24..bf97473544 100644 --- a/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/inbound/KafkaInboundGateway.java +++ b/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/inbound/KafkaInboundGateway.java @@ -41,7 +41,6 @@ import org.springframework.kafka.listener.ConsumerSeekAware; import org.springframework.kafka.listener.ContainerProperties; import org.springframework.kafka.listener.MessageListener; import org.springframework.kafka.listener.adapter.RecordMessagingMessageListenerAdapter; -import org.springframework.kafka.listener.adapter.RetryingMessageListenerAdapter; import org.springframework.kafka.support.Acknowledgment; import org.springframework.kafka.support.KafkaHeaders; import org.springframework.kafka.support.converter.ConversionException; @@ -130,7 +129,7 @@ public class KafkaInboundGateway extends MessagingGatewaySupport implem /** * Specify a {@link RetryTemplate} instance to wrap * {@link KafkaInboundGateway.IntegrationRecordMessageListener} into - * {@link RetryingMessageListenerAdapter}. + * {@code RetryingMessageListenerAdapter}. * @param retryTemplate the {@link RetryTemplate} to use. */ public void setRetryTemplate(RetryTemplate retryTemplate) { @@ -174,12 +173,14 @@ public class KafkaInboundGateway extends MessagingGatewaySupport implem } @Override + @SuppressWarnings("deprecation") protected void onInit() { super.onInit(); MessageListener kafkaListener = this.listener; if (this.retryTemplate != null) { kafkaListener = - new RetryingMessageListenerAdapter<>(kafkaListener, this.retryTemplate, this.recoveryCallback); + new org.springframework.kafka.listener.adapter.RetryingMessageListenerAdapter<>(kafkaListener, + this.retryTemplate, this.recoveryCallback); this.retryTemplate.registerListener(this.listener); } ContainerProperties containerProperties = this.messageListenerContainer.getContainerProperties(); diff --git a/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/inbound/KafkaMessageDrivenChannelAdapter.java b/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/inbound/KafkaMessageDrivenChannelAdapter.java index c25ac04eb3..45f9401004 100644 --- a/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/inbound/KafkaMessageDrivenChannelAdapter.java +++ b/spring-integration-kafka/src/main/java/org/springframework/integration/kafka/inbound/KafkaMessageDrivenChannelAdapter.java @@ -46,7 +46,6 @@ import org.springframework.kafka.listener.adapter.FilteringBatchMessageListenerA import org.springframework.kafka.listener.adapter.FilteringMessageListenerAdapter; import org.springframework.kafka.listener.adapter.RecordFilterStrategy; import org.springframework.kafka.listener.adapter.RecordMessagingMessageListenerAdapter; -import org.springframework.kafka.listener.adapter.RetryingMessageListenerAdapter; import org.springframework.kafka.support.Acknowledgment; import org.springframework.kafka.support.KafkaHeaders; import org.springframework.kafka.support.converter.BatchMessageConverter; @@ -194,7 +193,7 @@ public class KafkaMessageDrivenChannelAdapter extends MessageProducerSuppo /** * Specify a {@link RetryTemplate} instance to wrap * {@link KafkaMessageDrivenChannelAdapter.IntegrationRecordMessageListener} into - * {@link RetryingMessageListenerAdapter}. + * {@code RetryingMessageListenerAdapter}. * @param retryTemplate the {@link RetryTemplate} to use. * @since 2.0.1 */ @@ -218,12 +217,12 @@ public class KafkaMessageDrivenChannelAdapter extends MessageProducerSuppo /** * The {@code boolean} flag to specify the order how - * {@link RetryingMessageListenerAdapter} and + * {@code RetryingMessageListenerAdapter} and * {@link FilteringMessageListenerAdapter} are wrapped to each other, * if both of them are present. * Does not make sense if only one of {@link RetryTemplate} or * {@link RecordFilterStrategy} is present, or any. - * @param filterInRetry the order for {@link RetryingMessageListenerAdapter} and + * @param filterInRetry the order for {@code RetryingMessageListenerAdapter} and * {@link FilteringMessageListenerAdapter} wrapping. Defaults to {@code false}. * @since 2.0.1 */ @@ -274,6 +273,7 @@ public class KafkaMessageDrivenChannelAdapter extends MessageProducerSuppo } @Override + @SuppressWarnings("deprecation") protected void onInit() { super.onInit(); @@ -292,14 +292,14 @@ public class KafkaMessageDrivenChannelAdapter extends MessageProducerSuppo if (doFilterInRetry) { listener = new FilteringMessageListenerAdapter<>(listener, this.recordFilterStrategy, this.ackDiscarded); - listener = new RetryingMessageListenerAdapter<>(listener, this.retryTemplate, - this.recoveryCallback); + listener = new org.springframework.kafka.listener.adapter.RetryingMessageListenerAdapter<>(listener, + this.retryTemplate, this.recoveryCallback); this.retryTemplate.registerListener(this.recordListener); } else { if (this.retryTemplate != null) { - listener = new RetryingMessageListenerAdapter<>(listener, this.retryTemplate, - this.recoveryCallback); + listener = new org.springframework.kafka.listener.adapter.RetryingMessageListenerAdapter<>(listener, + this.retryTemplate, this.recoveryCallback); this.retryTemplate.registerListener(this.recordListener); } if (this.recordFilterStrategy != null) {