GH-124: Add Retry and Filter Option to MDAdapter
Fixes GH-124 (https://github.com/spring-projects/spring-integration-kafka/issues/124) * Add `RecordFilterStrategy`, `ackDiscarded`, `retryTemplate`, `RecoveryCallback` and `filterInRetry` options into `KafkaMessageDrivenChannelAdapter` * Wrap an internal `IntegrationMessageListener` into `FilteringAcknowledgingMessageListenerAdapter` and/or `RetryingAcknowledgingMessageListenerAdapter` if those options are provided. * The `filterInRetry` flag provides a logic to identify the wrapping order for those adapters. Change `ContainerProperties` instance usage in the test-case from original to an internal in the `KafkaMessageListenerContainer` That meets changes for the https://github.com/spring-projects/spring-kafka/issues/154
This commit is contained in:
@@ -21,10 +21,16 @@ import org.apache.kafka.clients.consumer.ConsumerRecord;
|
||||
import org.springframework.integration.context.OrderlyShutdownCapable;
|
||||
import org.springframework.integration.endpoint.MessageProducerSupport;
|
||||
import org.springframework.kafka.listener.AbstractMessageListenerContainer;
|
||||
import org.springframework.kafka.listener.AcknowledgingMessageListener;
|
||||
import org.springframework.kafka.listener.adapter.FilteringAcknowledgingMessageListenerAdapter;
|
||||
import org.springframework.kafka.listener.adapter.MessagingMessageListenerAdapter;
|
||||
import org.springframework.kafka.listener.adapter.RecordFilterStrategy;
|
||||
import org.springframework.kafka.listener.adapter.RetryingAcknowledgingMessageListenerAdapter;
|
||||
import org.springframework.kafka.support.Acknowledgment;
|
||||
import org.springframework.kafka.support.converter.MessageConverter;
|
||||
import org.springframework.messaging.Message;
|
||||
import org.springframework.retry.RecoveryCallback;
|
||||
import org.springframework.retry.support.RetryTemplate;
|
||||
import org.springframework.util.Assert;
|
||||
|
||||
/**
|
||||
@@ -44,19 +50,115 @@ public class KafkaMessageDrivenChannelAdapter<K, V> extends MessageProducerSuppo
|
||||
|
||||
private final MessagingMessageListenerAdapter<K, V> listener = new IntegrationMessageListener();
|
||||
|
||||
private RecordFilterStrategy<K, V> recordFilterStrategy;
|
||||
|
||||
private boolean ackDiscarded;
|
||||
|
||||
private RetryTemplate retryTemplate;
|
||||
|
||||
private RecoveryCallback<Void> recoveryCallback;
|
||||
|
||||
private boolean filterInRetry;
|
||||
|
||||
public KafkaMessageDrivenChannelAdapter(AbstractMessageListenerContainer<K, V> messageListenerContainer) {
|
||||
Assert.notNull(messageListenerContainer, "messageListenerContainer is required");
|
||||
Assert.isNull(messageListenerContainer.getContainerProperties().getMessageListener(),
|
||||
"Container must not already have a listener");
|
||||
this.messageListenerContainer = messageListenerContainer;
|
||||
this.messageListenerContainer.setAutoStartup(false);
|
||||
this.messageListenerContainer.getContainerProperties().setMessageListener(this.listener);
|
||||
}
|
||||
|
||||
public void setMessageConverter(MessageConverter messageConverter) {
|
||||
this.listener.setMessageConverter(messageConverter);
|
||||
}
|
||||
|
||||
/**
|
||||
* Specify a {@link RecordFilterStrategy} to wrap
|
||||
* {@link KafkaMessageDrivenChannelAdapter.IntegrationMessageListener} into
|
||||
* {@link FilteringAcknowledgingMessageListenerAdapter}.
|
||||
* @param recordFilterStrategy the {@link RecordFilterStrategy} to use.
|
||||
* @since 2.0.1
|
||||
*/
|
||||
public void setRecordFilterStrategy(RecordFilterStrategy<K, V> recordFilterStrategy) {
|
||||
this.recordFilterStrategy = recordFilterStrategy;
|
||||
}
|
||||
|
||||
/**
|
||||
* A {@code boolean} flag to indicate if {@link FilteringAcknowledgingMessageListenerAdapter}
|
||||
* should acknowledge discarded records or not.
|
||||
* Does not make sense if {@link #setRecordFilterStrategy(RecordFilterStrategy)} isn't specified.
|
||||
* @param ackDiscarded true to ack (commit offset for) discarded messages.
|
||||
* @since 2.0.1
|
||||
*/
|
||||
public void setAckDiscarded(boolean ackDiscarded) {
|
||||
this.ackDiscarded = ackDiscarded;
|
||||
}
|
||||
|
||||
/**
|
||||
* Specify a {@link RetryTemplate} instance to wrap
|
||||
* {@link KafkaMessageDrivenChannelAdapter.IntegrationMessageListener} into
|
||||
* {@link RetryingAcknowledgingMessageListenerAdapter}.
|
||||
* @param retryTemplate the {@link RetryTemplate} to use.
|
||||
* @since 2.0.1
|
||||
*/
|
||||
public void setRetryTemplate(RetryTemplate retryTemplate) {
|
||||
this.retryTemplate = retryTemplate;
|
||||
}
|
||||
|
||||
/**
|
||||
* A {@link RecoveryCallback} instance for retry operation;
|
||||
* if null, the exception will be thrown to the container after retries are exhausted.
|
||||
* Does not make sense if {@link #setRetryTemplate(RetryTemplate)} isn't specified.
|
||||
* @param recoveryCallback the recovery callback.
|
||||
* @since 2.0.1
|
||||
*/
|
||||
public void setRecoveryCallback(RecoveryCallback<Void> recoveryCallback) {
|
||||
this.recoveryCallback = recoveryCallback;
|
||||
}
|
||||
|
||||
/**
|
||||
* The {@code boolean} flag to specify the order how
|
||||
* {@link RetryingAcknowledgingMessageListenerAdapter} and
|
||||
* {@link FilteringAcknowledgingMessageListenerAdapter} 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 RetryingAcknowledgingMessageListenerAdapter} and
|
||||
* {@link FilteringAcknowledgingMessageListenerAdapter} wrapping. Defaults to {@code false}.
|
||||
* @since 2.0.1
|
||||
*/
|
||||
public void setFilterInRetry(boolean filterInRetry) {
|
||||
this.filterInRetry = filterInRetry;
|
||||
}
|
||||
|
||||
@Override
|
||||
protected void onInit() {
|
||||
super.onInit();
|
||||
|
||||
AcknowledgingMessageListener<K, V> listener = this.listener;
|
||||
|
||||
boolean filterInRetry = this.filterInRetry && this.retryTemplate != null && this.recordFilterStrategy != null;
|
||||
|
||||
if (filterInRetry) {
|
||||
listener = new FilteringAcknowledgingMessageListenerAdapter<>(listener, this.recordFilterStrategy,
|
||||
this.ackDiscarded);
|
||||
listener = new RetryingAcknowledgingMessageListenerAdapter<>(listener, this.retryTemplate,
|
||||
this.recoveryCallback);
|
||||
}
|
||||
else {
|
||||
if (this.retryTemplate != null) {
|
||||
listener = new RetryingAcknowledgingMessageListenerAdapter<>(listener, this.retryTemplate,
|
||||
this.recoveryCallback);
|
||||
}
|
||||
if (this.recordFilterStrategy != null) {
|
||||
listener = new FilteringAcknowledgingMessageListenerAdapter<>(listener, this.recordFilterStrategy,
|
||||
this.ackDiscarded);
|
||||
}
|
||||
}
|
||||
|
||||
this.messageListenerContainer.getContainerProperties().setMessageListener(listener);
|
||||
}
|
||||
|
||||
@Override
|
||||
protected void doStart() {
|
||||
this.messageListenerContainer.start();
|
||||
|
||||
@@ -17,6 +17,9 @@
|
||||
package org.springframework.integration.kafka.config.xml;
|
||||
|
||||
import static org.assertj.core.api.Assertions.assertThat;
|
||||
import static org.mockito.Mockito.mock;
|
||||
|
||||
import java.util.Collections;
|
||||
|
||||
import org.junit.Test;
|
||||
import org.junit.runner.RunWith;
|
||||
@@ -24,9 +27,16 @@ import org.junit.runner.RunWith;
|
||||
import org.springframework.beans.factory.annotation.Autowired;
|
||||
import org.springframework.integration.channel.NullChannel;
|
||||
import org.springframework.integration.channel.PublishSubscribeChannel;
|
||||
import org.springframework.integration.channel.QueueChannel;
|
||||
import org.springframework.integration.kafka.inbound.KafkaMessageDrivenChannelAdapter;
|
||||
import org.springframework.integration.test.util.TestUtils;
|
||||
import org.springframework.kafka.core.DefaultKafkaConsumerFactory;
|
||||
import org.springframework.kafka.listener.KafkaMessageListenerContainer;
|
||||
import org.springframework.kafka.listener.adapter.FilteringAcknowledgingMessageListenerAdapter;
|
||||
import org.springframework.kafka.listener.adapter.RecordFilterStrategy;
|
||||
import org.springframework.kafka.listener.adapter.RetryingAcknowledgingMessageListenerAdapter;
|
||||
import org.springframework.kafka.listener.config.ContainerProperties;
|
||||
import org.springframework.retry.support.RetryTemplate;
|
||||
import org.springframework.test.context.ContextConfiguration;
|
||||
import org.springframework.test.context.junit4.SpringJUnit4ClassRunner;
|
||||
|
||||
@@ -60,4 +70,59 @@ public class KafkaMessageDrivenChannelAdapterParserTests {
|
||||
assertThat(container).isNotNull();
|
||||
}
|
||||
|
||||
@Test
|
||||
@SuppressWarnings("unchecked")
|
||||
public void testKafkaMessageDrivenChannelAdapterOptions() {
|
||||
DefaultKafkaConsumerFactory<Integer, String> cf =
|
||||
new DefaultKafkaConsumerFactory<>(Collections.<String, Object>emptyMap());
|
||||
ContainerProperties containerProps = new ContainerProperties("foo");
|
||||
KafkaMessageListenerContainer<Integer, String> container =
|
||||
new KafkaMessageListenerContainer<>(cf, containerProps);
|
||||
KafkaMessageDrivenChannelAdapter<Integer, String> adapter = new KafkaMessageDrivenChannelAdapter<>(container);
|
||||
adapter.setOutputChannel(new QueueChannel());
|
||||
|
||||
adapter.setRecordFilterStrategy(mock(RecordFilterStrategy.class));
|
||||
adapter.afterPropertiesSet();
|
||||
|
||||
containerProps = TestUtils.getPropertyValue(container, "containerProperties", ContainerProperties.class);
|
||||
|
||||
Object messageListener = containerProps.getMessageListener();
|
||||
assertThat(messageListener).isInstanceOf(FilteringAcknowledgingMessageListenerAdapter.class);
|
||||
|
||||
Object delegate = TestUtils.getPropertyValue(messageListener, "delegate");
|
||||
|
||||
assertThat(delegate.getClass().getName()).contains("$IntegrationMessageListener");
|
||||
|
||||
adapter.setRecordFilterStrategy(null);
|
||||
adapter.setRetryTemplate(new RetryTemplate());
|
||||
adapter.afterPropertiesSet();
|
||||
|
||||
messageListener = containerProps.getMessageListener();
|
||||
assertThat(messageListener).isInstanceOf(RetryingAcknowledgingMessageListenerAdapter.class);
|
||||
|
||||
delegate = TestUtils.getPropertyValue(messageListener, "delegate");
|
||||
|
||||
assertThat(delegate.getClass().getName()).contains("$IntegrationMessageListener");
|
||||
|
||||
adapter.setRecordFilterStrategy(mock(RecordFilterStrategy.class));
|
||||
adapter.afterPropertiesSet();
|
||||
|
||||
messageListener = containerProps.getMessageListener();
|
||||
assertThat(messageListener).isInstanceOf(FilteringAcknowledgingMessageListenerAdapter.class);
|
||||
|
||||
delegate = TestUtils.getPropertyValue(messageListener, "delegate");
|
||||
|
||||
assertThat(delegate).isInstanceOf(RetryingAcknowledgingMessageListenerAdapter.class);
|
||||
|
||||
adapter.setFilterInRetry(true);
|
||||
adapter.afterPropertiesSet();
|
||||
|
||||
messageListener = containerProps.getMessageListener();
|
||||
assertThat(messageListener).isInstanceOf(RetryingAcknowledgingMessageListenerAdapter.class);
|
||||
|
||||
delegate = TestUtils.getPropertyValue(messageListener, "delegate");
|
||||
|
||||
assertThat(delegate).isInstanceOf(FilteringAcknowledgingMessageListenerAdapter.class);
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user