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
This commit is contained in:
@@ -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<K, V, R> 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<K, V, R> extends MessagingGatewaySupport implem
|
||||
}
|
||||
|
||||
@Override
|
||||
@SuppressWarnings("deprecation")
|
||||
protected void onInit() {
|
||||
super.onInit();
|
||||
MessageListener<K, V> 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();
|
||||
|
||||
@@ -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<K, V> 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<K, V> 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<K, V> extends MessageProducerSuppo
|
||||
}
|
||||
|
||||
@Override
|
||||
@SuppressWarnings("deprecation")
|
||||
protected void onInit() {
|
||||
super.onInit();
|
||||
|
||||
@@ -292,14 +292,14 @@ public class KafkaMessageDrivenChannelAdapter<K, V> 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) {
|
||||
|
||||
Reference in New Issue
Block a user